workflow

package
v0.5.12 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 13, 2026 License: Apache-2.0 Imports: 8 Imported by: 0

Documentation

Overview

Package workflow 提供 Hexagon AI Agent 框架的工作流引擎

工作流引擎支持定义和执行复杂的多步骤任务,具有以下特性: - 顺序执行、并行执行、条件分支 - 暂停/恢复执行 - 持久化和恢复 - 错误处理和重试

基本用法:

wf := workflow.New("my-workflow").
    Add(step1).
    Add(step2).
    Parallel(step3, step4).
    Build()

result, err := wf.Run(ctx, input)

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type BaseStep

type BaseStep struct {
	// contains filtered or unexported fields
}

BaseStep 基础步骤实现

func NewStep

func NewStep(id, name string, fn StepFunc, opts ...BaseStepOption) *BaseStep

NewStep 创建基础步骤

func (*BaseStep) Dependencies

func (s *BaseStep) Dependencies() []string

Dependencies 返回依赖的步骤 ID

func (*BaseStep) Execute

func (s *BaseStep) Execute(ctx context.Context, input StepInput) (*StepOutput, error)

Execute 执行步骤

func (*BaseStep) ID

func (s *BaseStep) ID() string

ID 返回步骤 ID

func (*BaseStep) Name

func (s *BaseStep) Name() string

Name 返回步骤名称

func (*BaseStep) Type

func (s *BaseStep) Type() StepType

Type 返回步骤类型

func (*BaseStep) Validate

func (s *BaseStep) Validate() error

Validate 验证步骤配置

type BaseStepOption

type BaseStepOption func(*BaseStep)

BaseStepOption 基础步骤选项

func WithStepDependencies

func WithStepDependencies(deps ...string) BaseStepOption

WithStepDependencies 设置步骤依赖

func WithStepDescription

func WithStepDescription(desc string) BaseStepOption

WithStepDescription 设置步骤描述

func WithStepMetadata

func WithStepMetadata(key string, value any) BaseStepOption

WithStepMetadata 设置步骤元数据

func WithStepRetryPolicy

func WithStepRetryPolicy(policy *RetryPolicy) BaseStepOption

WithStepRetryPolicy 设置步骤重试策略

func WithStepTimeout

func WithStepTimeout(timeout time.Duration) BaseStepOption

WithStepTimeout 设置步骤超时时间

type ConditionFunc

type ConditionFunc func(ctx context.Context, input StepInput) (string, error)

ConditionFunc 条件函数

type ConditionalBuilder

type ConditionalBuilder struct {
	// contains filtered or unexported fields
}

ConditionalBuilder 条件构建器

func (*ConditionalBuilder) Branch

func (b *ConditionalBuilder) Branch(name string, step Step) *ConditionalBuilder

Branch 添加分支

func (*ConditionalBuilder) Else

Else 设置条件为假时执行的步骤

func (*ConditionalBuilder) ElseFunc

func (b *ConditionalBuilder) ElseFunc(id, name string, fn StepFunc) *ConditionalBuilder

ElseFunc 设置条件为假时执行的函数

func (*ConditionalBuilder) End

End 结束条件构建

func (*ConditionalBuilder) Then

Then 设置条件为真时执行的步骤

func (*ConditionalBuilder) ThenFunc

func (b *ConditionalBuilder) ThenFunc(id, name string, fn StepFunc) *ConditionalBuilder

ThenFunc 设置条件为真时执行的函数

type ConditionalStep

type ConditionalStep struct {
	// contains filtered or unexported fields
}

ConditionalStep 条件步骤

func NewConditionalStep

func NewConditionalStep(id, name string, condition ConditionFunc) *ConditionalStep

NewConditionalStep 创建条件步骤

func (*ConditionalStep) Branch

func (s *ConditionalStep) Branch(name string, step Step) *ConditionalStep

Branch 添加分支

func (*ConditionalStep) Else

func (s *ConditionalStep) Else(step Step) *ConditionalStep

Else 设置条件为假时执行的步骤

func (*ConditionalStep) Execute

func (s *ConditionalStep) Execute(ctx context.Context, input StepInput) (*StepOutput, error)

Execute 执行条件步骤

func (*ConditionalStep) ID

func (s *ConditionalStep) ID() string

ID 返回步骤 ID

func (*ConditionalStep) Name

func (s *ConditionalStep) Name() string

Name 返回步骤名称

func (*ConditionalStep) Then

func (s *ConditionalStep) Then(step Step) *ConditionalStep

Then 设置条件为真时执行的步骤

func (*ConditionalStep) Type

func (s *ConditionalStep) Type() StepType

Type 返回步骤类型

func (*ConditionalStep) Validate

func (s *ConditionalStep) Validate() error

Validate 验证步骤配置

type ExecutionContext

type ExecutionContext struct {
	// Variables 变量
	Variables map[string]any `json:"variables,omitempty"`

	// CurrentStepID 当前步骤 ID
	CurrentStepID string `json:"current_step_id,omitempty"`

	// CompletedSteps 已完成的步骤
	CompletedSteps []string `json:"completed_steps,omitempty"`

	// PendingSteps 待执行的步骤
	PendingSteps []string `json:"pending_steps,omitempty"`

	// SkippedSteps 跳过的步骤
	SkippedSteps []string `json:"skipped_steps,omitempty"`

	// Metadata 元数据
	Metadata map[string]any `json:"metadata,omitempty"`
}

ExecutionContext 执行上下文

type ExecutionRecovery

type ExecutionRecovery struct {
	// contains filtered or unexported fields
}

ExecutionRecovery 执行恢复器

func NewExecutionRecovery

func NewExecutionRecovery(store WorkflowStore, registry *WorkflowRegistry, executor *Executor) *ExecutionRecovery

NewExecutionRecovery 创建执行恢复器

func (*ExecutionRecovery) RecoverPending

func (r *ExecutionRecovery) RecoverPending(ctx context.Context) (int, error)

RecoverPending 恢复待处理的执行

type ExecutionSnapshot

type ExecutionSnapshot struct {
	// ExecutionID 执行实例 ID
	ExecutionID string `json:"execution_id"`

	// WorkflowID 工作流 ID
	WorkflowID string `json:"workflow_id"`

	// Status 状态
	Status WorkflowStatus `json:"status"`

	// Context 上下文
	Context *ExecutionContext `json:"context"`

	// StepResults 步骤结果
	StepResults map[string]*StepResult `json:"step_results"`

	// Input 输入
	Input json.RawMessage `json:"input"`

	// CreatedAt 创建时间
	CreatedAt time.Time `json:"created_at"`
}

ExecutionSnapshot 执行快照(用于恢复)

func CreateSnapshot

func CreateSnapshot(execution *WorkflowExecution) *ExecutionSnapshot

CreateSnapshot 创建执行快照

type Executor

type Executor struct {
	// contains filtered or unexported fields
}

Executor 工作流执行器

func NewExecutor

func NewExecutor(opts ...ExecutorOption) *Executor

NewExecutor 创建执行器

func (*Executor) Cancel

func (e *Executor) Cancel(ctx context.Context, executionID string) error

Cancel 取消执行

func (*Executor) CleanupCompleted

func (e *Executor) CleanupCompleted(olderThan time.Duration) int

CleanupCompleted 清理已完成的执行

func (*Executor) GetExecution

func (e *Executor) GetExecution(ctx context.Context, executionID string) (*WorkflowExecution, error)

GetExecution 获取执行实例

func (*Executor) ListExecutions

func (e *Executor) ListExecutions(ctx context.Context, workflowID string, status WorkflowStatus, limit int) ([]*WorkflowExecution, error)

ListExecutions 列出执行实例

func (*Executor) OnEvent

func (e *Executor) OnEvent(handler WorkflowEventHandler)

OnEvent 注册事件处理器

func (*Executor) Pause

func (e *Executor) Pause(ctx context.Context, executionID string) error

Pause 暂停执行

func (*Executor) Resume

func (e *Executor) Resume(ctx context.Context, executionID string) error

Resume 恢复执行

func (*Executor) Run

func (e *Executor) Run(ctx context.Context, wf *Workflow, input WorkflowInput) (*WorkflowOutput, error)

Run 同步运行工作流

func (*Executor) RunAsync

func (e *Executor) RunAsync(ctx context.Context, wf *Workflow, input WorkflowInput) (string, error)

RunAsync 异步运行工作流

func (*Executor) WaitForCompletion

func (e *Executor) WaitForCompletion(ctx context.Context, executionID string, timeout time.Duration) (*WorkflowExecution, error)

WaitForCompletion 等待执行完成

type ExecutorConfig

type ExecutorConfig struct {
	// DefaultTimeout 默认超时时间
	DefaultTimeout time.Duration

	// MaxConcurrentExecutions 最大并发执行数
	MaxConcurrentExecutions int

	// EnablePersistence 启用持久化
	EnablePersistence bool
}

ExecutorConfig 执行器配置

func DefaultExecutorConfig

func DefaultExecutorConfig() ExecutorConfig

DefaultExecutorConfig 返回默认配置

type ExecutorOption

type ExecutorOption func(*Executor)

ExecutorOption 执行器选项

func WithExecutorConfig

func WithExecutorConfig(config ExecutorConfig) ExecutorOption

WithExecutorConfig 设置配置

func WithHooks

func WithHooks(hooks *WorkflowHooks) ExecutorOption

WithHooks 设置钩子

func WithStore

func WithStore(store WorkflowStore) ExecutorOption

WithStore 设置存储

type LoopStep

type LoopStep struct {
	// contains filtered or unexported fields
}

LoopStep 循环步骤

func NewLoopStep

func NewLoopStep(id, name string, step Step, condition ConditionFunc, opts ...LoopStepOption) *LoopStep

NewLoopStep 创建循环步骤 condition 返回 "continue" 继续循环,返回 "break" 退出

func (*LoopStep) Execute

func (s *LoopStep) Execute(ctx context.Context, input StepInput) (*StepOutput, error)

Execute 执行循环步骤

func (*LoopStep) ID

func (s *LoopStep) ID() string

ID 返回步骤 ID

func (*LoopStep) Name

func (s *LoopStep) Name() string

Name 返回步骤名称

func (*LoopStep) Type

func (s *LoopStep) Type() StepType

Type 返回步骤类型

func (*LoopStep) Validate

func (s *LoopStep) Validate() error

Validate 验证步骤配置

type LoopStepOption

type LoopStepOption func(*LoopStep)

LoopStepOption 循环步骤选项

func WithCollectOutput

func WithCollectOutput(collect bool) LoopStepOption

WithCollectOutput 设置是否收集输出

func WithMaxIterations

func WithMaxIterations(max int) LoopStepOption

WithMaxIterations 设置最大迭代次数

type MemoryWorkflowStore

type MemoryWorkflowStore struct {
	// contains filtered or unexported fields
}

MemoryWorkflowStore 内存存储实现

func NewMemoryWorkflowStore

func NewMemoryWorkflowStore() *MemoryWorkflowStore

NewMemoryWorkflowStore 创建内存存储

func (*MemoryWorkflowStore) Clear

func (s *MemoryWorkflowStore) Clear()

Clear 清空存储

func (*MemoryWorkflowStore) DeleteExecution

func (s *MemoryWorkflowStore) DeleteExecution(ctx context.Context, id string) error

DeleteExecution 删除执行实例

func (*MemoryWorkflowStore) DeleteWorkflow

func (s *MemoryWorkflowStore) DeleteWorkflow(ctx context.Context, id string) error

DeleteWorkflow 删除工作流定义

func (*MemoryWorkflowStore) GetExecution

func (s *MemoryWorkflowStore) GetExecution(ctx context.Context, id string) (*WorkflowExecution, error)

GetExecution 获取执行实例

func (*MemoryWorkflowStore) GetPendingExecutions

func (s *MemoryWorkflowStore) GetPendingExecutions(ctx context.Context, limit int) ([]*WorkflowExecution, error)

GetPendingExecutions 获取待恢复的执行实例

func (*MemoryWorkflowStore) GetWorkflow

func (s *MemoryWorkflowStore) GetWorkflow(ctx context.Context, id string) (*Workflow, error)

GetWorkflow 获取工作流定义

func (*MemoryWorkflowStore) ListExecutions

func (s *MemoryWorkflowStore) ListExecutions(ctx context.Context, workflowID string, status WorkflowStatus, limit int) ([]*WorkflowExecution, error)

ListExecutions 列出执行实例

func (*MemoryWorkflowStore) ListWorkflows

func (s *MemoryWorkflowStore) ListWorkflows(ctx context.Context, limit int) ([]*Workflow, error)

ListWorkflows 列出工作流定义

func (*MemoryWorkflowStore) SaveExecution

func (s *MemoryWorkflowStore) SaveExecution(ctx context.Context, execution *WorkflowExecution) error

SaveExecution 保存执行实例

func (*MemoryWorkflowStore) SaveWorkflow

func (s *MemoryWorkflowStore) SaveWorkflow(ctx context.Context, wf *Workflow) error

SaveWorkflow 保存工作流定义

func (*MemoryWorkflowStore) Stats

func (s *MemoryWorkflowStore) Stats() (workflowCount, executionCount int)

Stats 返回存储统计

func (*MemoryWorkflowStore) UpdateExecutionStatus

func (s *MemoryWorkflowStore) UpdateExecutionStatus(ctx context.Context, id string, status WorkflowStatus, errMsg string) error

UpdateExecutionStatus 更新执行状态

type ParallelStep

type ParallelStep struct {
	// contains filtered or unexported fields
}

ParallelStep 并行步骤

func NewParallelStep

func NewParallelStep(id, name string, steps []Step, opts ...ParallelStepOption) *ParallelStep

NewParallelStep 创建并行步骤

func (*ParallelStep) Execute

func (s *ParallelStep) Execute(ctx context.Context, input StepInput) (*StepOutput, error)

Execute 执行并行步骤

func (*ParallelStep) ID

func (s *ParallelStep) ID() string

ID 返回步骤 ID

func (*ParallelStep) Name

func (s *ParallelStep) Name() string

Name 返回步骤名称

func (*ParallelStep) Type

func (s *ParallelStep) Type() StepType

Type 返回步骤类型

func (*ParallelStep) Validate

func (s *ParallelStep) Validate() error

Validate 验证步骤配置

type ParallelStepOption

type ParallelStepOption func(*ParallelStep)

ParallelStepOption 并行步骤选项

func WithFailFast

func WithFailFast(failFast bool) ParallelStepOption

WithFailFast 设置快速失败

func WithMaxParallel

func WithMaxParallel(max int) ParallelStepOption

WithMaxParallel 设置最大并行数

type RetryPolicy

type RetryPolicy struct {
	// MaxRetries 最大重试次数
	MaxRetries int `json:"max_retries"`

	// InitialInterval 初始间隔
	InitialInterval time.Duration `json:"initial_interval"`

	// MaxInterval 最大间隔
	MaxInterval time.Duration `json:"max_interval"`

	// Multiplier 间隔倍数
	Multiplier float64 `json:"multiplier"`

	// RetryableErrors 可重试的错误类型
	RetryableErrors []string `json:"retryable_errors,omitempty"`
}

RetryPolicy 重试策略

func DefaultRetryPolicy

func DefaultRetryPolicy() *RetryPolicy

DefaultRetryPolicy 返回默认重试策略

type Step

type Step interface {
	// ID 返回步骤 ID
	ID() string

	// Name 返回步骤名称
	Name() string

	// Type 返回步骤类型
	Type() StepType

	// Execute 执行步骤
	Execute(ctx context.Context, input StepInput) (*StepOutput, error)

	// Validate 验证步骤配置
	Validate() error
}

Step 步骤接口

func FilterStep

func FilterStep(id, name string, predicate func(item any) bool) Step

FilterStep 创建 Filter 步骤(过滤数组元素)

func MapStep

func MapStep(id, name string, fn func(ctx context.Context, item any) (any, error)) Step

MapStep 创建 Map 步骤(对数组中的每个元素执行操作)

func ReduceStep

func ReduceStep(id, name string, initial any, reducer func(acc, item any) any) Step

ReduceStep 创建 Reduce 步骤(聚合数组元素)

func RetryWrapper

func RetryWrapper(step Step, policy *RetryPolicy) Step

RetryWrapper 包装步骤添加重试逻辑

func TimeoutWrapper

func TimeoutWrapper(step Step, timeout time.Duration) Step

TimeoutWrapper 包装步骤添加超时

type StepDefinition

type StepDefinition struct {
	// ID 步骤 ID
	ID string `json:"id"`

	// Name 步骤名称
	Name string `json:"name"`

	// Type 步骤类型
	Type StepType `json:"type"`

	// Description 描述
	Description string `json:"description,omitempty"`

	// Config 配置
	Config map[string]any `json:"config,omitempty"`

	// Dependencies 依赖的步骤 ID
	Dependencies []string `json:"dependencies,omitempty"`
}

StepDefinition 步骤定义(用于序列化)

type StepFunc

type StepFunc func(ctx context.Context, input StepInput) (*StepOutput, error)

StepFunc 步骤执行函数

type StepInput

type StepInput struct {
	// Data 输入数据
	Data any

	// Variables 上下文变量
	Variables map[string]any

	// PreviousOutputs 前置步骤的输出
	PreviousOutputs map[string]any

	// Metadata 元数据
	Metadata map[string]any
}

StepInput 步骤输入

type StepOutput

type StepOutput struct {
	// Data 输出数据
	Data any

	// Variables 更新的变量
	Variables map[string]any

	// Metadata 元数据
	Metadata map[string]any

	// NextStepID 下一步骤 ID(用于条件步骤)
	NextStepID string
}

StepOutput 步骤输出

type StepResult

type StepResult struct {
	// StepID 步骤 ID
	StepID string `json:"step_id"`

	// Status 状态
	Status WorkflowStatus `json:"status"`

	// Output 输出数据
	Output any `json:"output,omitempty"`

	// Error 错误信息
	Error string `json:"error,omitempty"`

	// StartedAt 开始时间
	StartedAt time.Time `json:"started_at"`

	// CompletedAt 完成时间
	CompletedAt *time.Time `json:"completed_at,omitempty"`

	// Duration 执行时长
	Duration time.Duration `json:"duration,omitempty"`

	// RetryCount 重试次数
	RetryCount int `json:"retry_count,omitempty"`
}

StepResult 步骤执行结果

type StepType

type StepType string

StepType 步骤类型

const (
	// StepTypeNormal 普通步骤
	StepTypeNormal StepType = "normal"
	// StepTypeParallel 并行步骤
	StepTypeParallel StepType = "parallel"
	// StepTypeConditional 条件步骤
	StepTypeConditional StepType = "conditional"
	// StepTypeLoop 循环步骤
	StepTypeLoop StepType = "loop"
	// StepTypeSubWorkflow 子工作流步骤
	StepTypeSubWorkflow StepType = "sub_workflow"
	// StepTypeWait 等待步骤
	StepTypeWait StepType = "wait"
)

type SubWorkflowStep

type SubWorkflowStep struct {
	// contains filtered or unexported fields
}

SubWorkflowStep 子工作流步骤

func NewSubWorkflowStep

func NewSubWorkflowStep(id, name string, workflow *Workflow, runner WorkflowRunner) *SubWorkflowStep

NewSubWorkflowStep 创建子工作流步骤

func (*SubWorkflowStep) Execute

func (s *SubWorkflowStep) Execute(ctx context.Context, input StepInput) (*StepOutput, error)

Execute 执行子工作流步骤

func (*SubWorkflowStep) ID

func (s *SubWorkflowStep) ID() string

ID 返回步骤 ID

func (*SubWorkflowStep) Name

func (s *SubWorkflowStep) Name() string

Name 返回步骤名称

func (*SubWorkflowStep) Type

func (s *SubWorkflowStep) Type() StepType

Type 返回步骤类型

func (*SubWorkflowStep) Validate

func (s *SubWorkflowStep) Validate() error

Validate 验证步骤配置

type WaitStep

type WaitStep struct {
	// contains filtered or unexported fields
}

WaitStep 等待步骤

func NewWaitStep

func NewWaitStep(id, name string, duration time.Duration) *WaitStep

NewWaitStep 创建等待步骤(固定时间)

func NewWaitUntilStep

func NewWaitUntilStep(id, name string, until func(ctx context.Context, input StepInput) (bool, error)) *WaitStep

NewWaitUntilStep 创建等待步骤(等待条件满足)

func (*WaitStep) Execute

func (s *WaitStep) Execute(ctx context.Context, input StepInput) (*StepOutput, error)

Execute 执行等待步骤

func (*WaitStep) ID

func (s *WaitStep) ID() string

ID 返回步骤 ID

func (*WaitStep) Name

func (s *WaitStep) Name() string

Name 返回步骤名称

func (*WaitStep) Type

func (s *WaitStep) Type() StepType

Type 返回步骤类型

func (*WaitStep) Validate

func (s *WaitStep) Validate() error

Validate 验证步骤配置

type Workflow

type Workflow struct {
	// ID 唯一标识符
	ID string `json:"id"`

	// Name 工作流名称
	Name string `json:"name"`

	// Description 描述
	Description string `json:"description,omitempty"`

	// Version 版本
	Version string `json:"version,omitempty"`

	// Steps 步骤列表
	Steps []Step `json:"-"`

	// StepDefs 步骤定义(用于序列化)
	StepDefs []StepDefinition `json:"steps"`

	// Metadata 元数据
	Metadata map[string]any `json:"metadata,omitempty"`

	// Timeout 总超时时间
	Timeout time.Duration `json:"timeout,omitempty"`

	// RetryPolicy 默认重试策略
	RetryPolicy *RetryPolicy `json:"retry_policy,omitempty"`

	// CreatedAt 创建时间
	CreatedAt time.Time `json:"created_at"`
}

Workflow 工作流定义

func Pipeline

func Pipeline(name string, steps ...Step) (*Workflow, error)

Pipeline 创建流水线(简化的顺序工作流)

func PipelineFuncs

func PipelineFuncs(name string, funcs []struct {
	ID   string
	Name string
	Fn   StepFunc
}) (*Workflow, error)

PipelineFuncs 从函数列表创建流水线

type WorkflowBuilder

type WorkflowBuilder struct {
	// contains filtered or unexported fields
}

WorkflowBuilder 工作流构建器

func New

func New(name string) *WorkflowBuilder

New 创建工作流构建器

func (*WorkflowBuilder) Add

func (b *WorkflowBuilder) Add(step Step) *WorkflowBuilder

Add 添加步骤

func (*WorkflowBuilder) AddFunc

func (b *WorkflowBuilder) AddFunc(id, name string, fn StepFunc, opts ...BaseStepOption) *WorkflowBuilder

AddFunc 添加函数步骤

func (*WorkflowBuilder) Build

func (b *WorkflowBuilder) Build() (*Workflow, error)

Build 构建工作流

func (*WorkflowBuilder) Conditional

func (b *WorkflowBuilder) Conditional(id, name string, condition ConditionFunc) *ConditionalBuilder

Conditional 条件分支

func (*WorkflowBuilder) Loop

func (b *WorkflowBuilder) Loop(id, name string, step Step, condition ConditionFunc, opts ...LoopStepOption) *WorkflowBuilder

Loop 循环执行

func (*WorkflowBuilder) MustBuild

func (b *WorkflowBuilder) MustBuild() *Workflow

MustBuild 构建工作流,失败时 panic

⚠️ 警告:构建失败时会 panic。 仅在初始化时使用,不要在运行时调用。 推荐使用 Build() 方法并正确处理错误。

使用场景:

  • 程序启动时的全局初始化
  • 测试代码中

func (*WorkflowBuilder) Parallel

func (b *WorkflowBuilder) Parallel(id, name string, steps ...Step) *WorkflowBuilder

Parallel 并行执行多个步骤

func (*WorkflowBuilder) ParallelFuncs

func (b *WorkflowBuilder) ParallelFuncs(id, name string, funcs map[string]StepFunc) *WorkflowBuilder

ParallelFuncs 并行执行多个函数

func (*WorkflowBuilder) Sequential

func (b *WorkflowBuilder) Sequential(steps ...Step) *WorkflowBuilder

Sequential 顺序执行多个步骤

func (*WorkflowBuilder) SubWorkflow

func (b *WorkflowBuilder) SubWorkflow(id, name string, workflow *Workflow, runner WorkflowRunner) *WorkflowBuilder

SubWorkflow 子工作流

func (*WorkflowBuilder) Wait

func (b *WorkflowBuilder) Wait(id, name string, duration time.Duration) *WorkflowBuilder

Wait 等待固定时间

func (*WorkflowBuilder) WaitUntil

func (b *WorkflowBuilder) WaitUntil(id, name string, condition func(ctx context.Context, input StepInput) (bool, error)) *WorkflowBuilder

WaitUntil 等待条件满足

func (*WorkflowBuilder) WithDescription

func (b *WorkflowBuilder) WithDescription(desc string) *WorkflowBuilder

WithDescription 设置描述

func (*WorkflowBuilder) WithID

func (b *WorkflowBuilder) WithID(id string) *WorkflowBuilder

WithID 设置工作流 ID

func (*WorkflowBuilder) WithMetadata

func (b *WorkflowBuilder) WithMetadata(key string, value any) *WorkflowBuilder

WithMetadata 设置元数据

func (*WorkflowBuilder) WithRetryPolicy

func (b *WorkflowBuilder) WithRetryPolicy(policy *RetryPolicy) *WorkflowBuilder

WithRetryPolicy 设置重试策略

func (*WorkflowBuilder) WithTimeout

func (b *WorkflowBuilder) WithTimeout(timeout time.Duration) *WorkflowBuilder

WithTimeout 设置超时时间

func (*WorkflowBuilder) WithVersion

func (b *WorkflowBuilder) WithVersion(version string) *WorkflowBuilder

WithVersion 设置版本

type WorkflowEvent

type WorkflowEvent struct {
	// Type 事件类型
	Type WorkflowEventType `json:"type"`

	// ExecutionID 执行实例 ID
	ExecutionID string `json:"execution_id"`

	// StepID 步骤 ID(如果适用)
	StepID string `json:"step_id,omitempty"`

	// Status 状态
	Status WorkflowStatus `json:"status"`

	// Data 事件数据
	Data any `json:"data,omitempty"`

	// Error 错误信息
	Error string `json:"error,omitempty"`

	// Timestamp 时间戳
	Timestamp time.Time `json:"timestamp"`
}

WorkflowEvent 工作流事件

type WorkflowEventHandler

type WorkflowEventHandler func(event *WorkflowEvent)

WorkflowEventHandler 工作流事件处理器

type WorkflowEventType

type WorkflowEventType string

WorkflowEventType 工作流事件类型

const (
	// EventWorkflowStarted 工作流开始
	EventWorkflowStarted WorkflowEventType = "workflow_started"
	// EventWorkflowCompleted 工作流完成
	EventWorkflowCompleted WorkflowEventType = "workflow_completed"
	// EventWorkflowFailed 工作流失败
	EventWorkflowFailed WorkflowEventType = "workflow_failed"
	// EventWorkflowPaused 工作流暂停
	EventWorkflowPaused WorkflowEventType = "workflow_paused"
	// EventWorkflowResumed 工作流恢复
	EventWorkflowResumed WorkflowEventType = "workflow_resumed"
	// EventWorkflowCancelled 工作流取消
	EventWorkflowCancelled WorkflowEventType = "workflow_cancelled"
	// EventStepStarted 步骤开始
	EventStepStarted WorkflowEventType = "step_started"
	// EventStepCompleted 步骤完成
	EventStepCompleted WorkflowEventType = "step_completed"
	// EventStepFailed 步骤失败
	EventStepFailed WorkflowEventType = "step_failed"
	// EventStepSkipped 步骤跳过
	EventStepSkipped WorkflowEventType = "step_skipped"
	// EventStepRetrying 步骤重试
	EventStepRetrying WorkflowEventType = "step_retrying"
)

type WorkflowExecution

type WorkflowExecution struct {
	// ID 执行实例 ID
	ID string `json:"id"`

	// WorkflowID 工作流 ID
	WorkflowID string `json:"workflow_id"`

	// Status 状态
	Status WorkflowStatus `json:"status"`

	// Input 输入数据
	Input json.RawMessage `json:"input,omitempty"`

	// Output 输出数据
	Output json.RawMessage `json:"output,omitempty"`

	// Error 错误信息
	Error string `json:"error,omitempty"`

	// Context 执行上下文
	Context *ExecutionContext `json:"context"`

	// StepResults 步骤执行结果
	StepResults map[string]*StepResult `json:"step_results"`

	// StartedAt 开始时间
	StartedAt time.Time `json:"started_at"`

	// CompletedAt 完成时间
	CompletedAt *time.Time `json:"completed_at,omitempty"`

	// PausedAt 暂停时间
	PausedAt *time.Time `json:"paused_at,omitempty"`

	// Duration 执行时长
	Duration time.Duration `json:"duration,omitempty"`
}

WorkflowExecution 工作流执行实例

func NewExecution

func NewExecution(wf *Workflow) *WorkflowExecution

NewExecution 创建执行实例

func RestoreFromSnapshot

func RestoreFromSnapshot(snapshot *ExecutionSnapshot) *WorkflowExecution

RestoreFromSnapshot 从快照恢复执行

type WorkflowHooks

type WorkflowHooks struct {
	// OnStart 工作流开始时触发
	OnStart func(ctx context.Context, wf *Workflow, input WorkflowInput) error

	// OnComplete 工作流完成时触发
	OnComplete func(ctx context.Context, wf *Workflow, output *WorkflowOutput) error

	// OnError 工作流出错时触发
	OnError func(ctx context.Context, wf *Workflow, err error) error

	// OnStepStart 步骤开始时触发
	OnStepStart func(ctx context.Context, step Step, input any) error

	// OnStepComplete 步骤完成时触发
	OnStepComplete func(ctx context.Context, step Step, output any) error

	// OnStepError 步骤出错时触发
	OnStepError func(ctx context.Context, step Step, err error) error
}

WorkflowHooks 工作流钩子

type WorkflowInput

type WorkflowInput struct {
	// Data 输入数据
	Data any `json:"data"`

	// Variables 初始变量
	Variables map[string]any `json:"variables,omitempty"`

	// Metadata 元数据
	Metadata map[string]any `json:"metadata,omitempty"`
}

WorkflowInput 工作流输入

type WorkflowOutput

type WorkflowOutput struct {
	// Data 输出数据
	Data any `json:"data"`

	// Variables 最终变量
	Variables map[string]any `json:"variables,omitempty"`

	// StepOutputs 各步骤输出
	StepOutputs map[string]any `json:"step_outputs,omitempty"`

	// Metadata 元数据
	Metadata map[string]any `json:"metadata,omitempty"`
}

WorkflowOutput 工作流输出

type WorkflowRegistry

type WorkflowRegistry struct {
	// contains filtered or unexported fields
}

WorkflowRegistry 工作流注册中心

func NewWorkflowRegistry

func NewWorkflowRegistry() *WorkflowRegistry

NewWorkflowRegistry 创建工作流注册中心

func (*WorkflowRegistry) Get

func (r *WorkflowRegistry) Get(id string) (*Workflow, bool)

Get 获取工作流

func (*WorkflowRegistry) List

func (r *WorkflowRegistry) List() []*Workflow

List 列出所有工作流

func (*WorkflowRegistry) Register

func (r *WorkflowRegistry) Register(wf *Workflow) error

Register 注册工作流

func (*WorkflowRegistry) Remove

func (r *WorkflowRegistry) Remove(id string)

Remove 移除工作流

type WorkflowRunner

type WorkflowRunner interface {
	// Run 运行工作流
	Run(ctx context.Context, wf *Workflow, input WorkflowInput) (*WorkflowOutput, error)

	// RunAsync 异步运行工作流
	RunAsync(ctx context.Context, wf *Workflow, input WorkflowInput) (string, error)

	// Pause 暂停执行
	Pause(ctx context.Context, executionID string) error

	// Resume 恢复执行
	Resume(ctx context.Context, executionID string) error

	// Cancel 取消执行
	Cancel(ctx context.Context, executionID string) error

	// GetExecution 获取执行实例
	GetExecution(ctx context.Context, executionID string) (*WorkflowExecution, error)

	// WaitForCompletion 等待执行完成
	WaitForCompletion(ctx context.Context, executionID string, timeout time.Duration) (*WorkflowExecution, error)
}

WorkflowRunner 工作流运行器接口

type WorkflowStatus

type WorkflowStatus string

WorkflowStatus 工作流状态

const (
	// StatusPending 等待执行
	StatusPending WorkflowStatus = "pending"
	// StatusRunning 执行中
	StatusRunning WorkflowStatus = "running"
	// StatusPaused 已暂停
	StatusPaused WorkflowStatus = "paused"
	// StatusCompleted 已完成
	StatusCompleted WorkflowStatus = "completed"
	// StatusFailed 执行失败
	StatusFailed WorkflowStatus = "failed"
	// StatusCancelled 已取消
	StatusCancelled WorkflowStatus = "cancelled"
)

type WorkflowStore

type WorkflowStore interface {
	// SaveWorkflow 保存工作流定义
	SaveWorkflow(ctx context.Context, wf *Workflow) error

	// GetWorkflow 获取工作流定义
	GetWorkflow(ctx context.Context, id string) (*Workflow, error)

	// ListWorkflows 列出工作流定义
	ListWorkflows(ctx context.Context, limit int) ([]*Workflow, error)

	// DeleteWorkflow 删除工作流定义
	DeleteWorkflow(ctx context.Context, id string) error

	// SaveExecution 保存执行实例
	SaveExecution(ctx context.Context, execution *WorkflowExecution) error

	// GetExecution 获取执行实例
	GetExecution(ctx context.Context, id string) (*WorkflowExecution, error)

	// ListExecutions 列出执行实例
	ListExecutions(ctx context.Context, workflowID string, status WorkflowStatus, limit int) ([]*WorkflowExecution, error)

	// DeleteExecution 删除执行实例
	DeleteExecution(ctx context.Context, id string) error

	// UpdateExecutionStatus 更新执行状态
	UpdateExecutionStatus(ctx context.Context, id string, status WorkflowStatus, errMsg string) error

	// GetPendingExecutions 获取待恢复的执行实例
	GetPendingExecutions(ctx context.Context, limit int) ([]*WorkflowExecution, error)
}

WorkflowStore 工作流存储接口

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL