Documentation
¶
Index ¶
- func BrokerMessageWorkflow(ctx workflow.Context, body []byte) error
- func WithClientHostPort(hostPort string) func(*ClientOptions)
- func WithClientNamespace(namespace string) func(*ClientOptions)
- type ClientOptions
- type ExecuteOptions
- type WorkerOptions
- type WorkflowClient
- func (wc *WorkflowClient) Cancel(ctx context.Context, workflowID, runID string) error
- func (wc *WorkflowClient) Close() error
- func (wc *WorkflowClient) Describe(ctx context.Context, workflowID, runID string) (*workflowservice.DescribeWorkflowExecutionResponse, error)
- func (wc *WorkflowClient) Execute(ctx context.Context, args any, opts ExecuteOptions) (string, error)
- func (wc *WorkflowClient) ExecuteSync(ctx context.Context, args any, opts ExecuteOptions) ([]byte, error)
- func (wc *WorkflowClient) NewWorker(opts WorkerOptions) (*WorkflowWorker, error)
- func (wc *WorkflowClient) Query(ctx context.Context, workflowID, runID, queryType string, arg any) (any, error)
- func (wc *WorkflowClient) Signal(ctx context.Context, workflowID, runID, signalName string, arg any) error
- func (wc *WorkflowClient) StartSimpleWorker(ctx context.Context, taskQueue string, ...) (*WorkflowWorker, error)
- func (wc *WorkflowClient) TemporalClient() client.Client
- func (wc *WorkflowClient) WithTracing()
- type WorkflowWorker
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func BrokerMessageWorkflow ¶
BrokerMessageWorkflow 是默认消息工作流。 它接收消息体,并委托已注册的 ProcessMessage Activity 处理。 复杂编排场景可通过 WorkerOptions.Workflows 注册自定义工作流。
func WithClientHostPort ¶
func WithClientHostPort(hostPort string) func(*ClientOptions)
WithClientHostPort 设置 Temporal Server 地址。
func WithClientNamespace ¶
func WithClientNamespace(namespace string) func(*ClientOptions)
WithClientNamespace 设置 Temporal 命名空间。
Types ¶
type ClientOptions ¶
type ExecuteOptions ¶
type ExecuteOptions struct {
// TaskQueue 是工作流使用的任务队列。
TaskQueue string
// WorkflowID 是工作流执行的唯一标识。
WorkflowID string
// WorkflowFn 是需要执行的工作流函数,空值时使用默认工作流。
WorkflowFn any
// RunTimeout 是单次工作流运行的最大耗时。
RunTimeout time.Duration
// ExecutionTimeout 是包含重试和 continue-as-new 在内的总执行超时。
ExecutionTimeout time.Duration
// TaskTimeout 是单个工作流任务超时。
TaskTimeout time.Duration
// RetryPolicy 是工作流重试策略。
RetryPolicy *temporal.RetryPolicy
// CronSchedule 是工作流定时调度表达式。
CronSchedule string
// IDReusePolicy 控制 WorkflowID 已存在时的复用行为。
IDReusePolicy enums.WorkflowIdReusePolicy
}
type WorkerOptions ¶
type WorkflowClient ¶
type WorkflowClient struct {
// contains filtered or unexported fields
}
WorkflowClient 提供 Temporal 工作流操作的高层封装。
func NewClient ¶
func NewClient(opts ...func(*ClientOptions)) (*WorkflowClient, error)
NewClient 创建并连接 Temporal 客户端。
func (*WorkflowClient) Cancel ¶
func (wc *WorkflowClient) Cancel(ctx context.Context, workflowID, runID string) error
Cancel 请求取消正在运行的工作流。
func (*WorkflowClient) Describe ¶
func (wc *WorkflowClient) Describe(ctx context.Context, workflowID, runID string) (*workflowservice.DescribeWorkflowExecutionResponse, error)
Describe 获取工作流执行描述信息。
func (*WorkflowClient) Execute ¶
func (wc *WorkflowClient) Execute(ctx context.Context, args any, opts ExecuteOptions) (string, error)
Execute 异步启动一次工作流执行。 返回本次执行的 run ID。
func (*WorkflowClient) ExecuteSync ¶
func (wc *WorkflowClient) ExecuteSync(ctx context.Context, args any, opts ExecuteOptions) ([]byte, error)
ExecuteSync 启动工作流并等待完成,返回工作流结果。
func (*WorkflowClient) NewWorker ¶
func (wc *WorkflowClient) NewWorker(opts WorkerOptions) (*WorkflowWorker, error)
NewWorker 创建 Temporal Worker。 创建后不会自动启动,需要调用 Start 开始轮询任务。
func (*WorkflowClient) Query ¶
func (wc *WorkflowClient) Query(ctx context.Context, workflowID, runID, queryType string, arg any) (any, error)
Query 查询正在运行的工作流状态。
func (*WorkflowClient) Signal ¶
func (wc *WorkflowClient) Signal(ctx context.Context, workflowID, runID, signalName string, arg any) error
Signal 向正在运行的工作流发送信号。
func (*WorkflowClient) StartSimpleWorker ¶
func (wc *WorkflowClient) StartSimpleWorker(ctx context.Context, taskQueue string, handler func(ctx context.Context, body []byte) error, opts ...func(*WorkerOptions)) (*WorkflowWorker, error)
StartSimpleWorker 创建带单一消息处理函数的 Worker。 这是从任务队列开始处理消息的最简方式。
func (*WorkflowClient) TemporalClient ¶
func (wc *WorkflowClient) TemporalClient() client.Client
TemporalClient 返回底层 Temporal SDK 客户端,供高级场景使用。
func (*WorkflowClient) WithTracing ¶
func (wc *WorkflowClient) WithTracing()
WithTracing 为客户端启用 OpenTelemetry 链路追踪。
type WorkflowWorker ¶
WorkflowWorker 管理负责轮询任务的 Temporal Worker。
func (*WorkflowWorker) IsRunning ¶
func (ww *WorkflowWorker) IsRunning() bool
IsRunning 返回 Worker 是否已经成功启动且尚未停止。
func (*WorkflowWorker) RegisterActivity ¶
func (ww *WorkflowWorker) RegisterActivity(fn any)
RegisterActivity 注册 Activity 函数或结构体。 必须在 Start 之前调用。
func (*WorkflowWorker) RegisterWorkflow ¶
func (ww *WorkflowWorker) RegisterWorkflow(fn any)
RegisterWorkflow 注册工作流函数。 必须在 Start 之前调用。
func (*WorkflowWorker) TaskQueue ¶
func (ww *WorkflowWorker) TaskQueue() string
TaskQueue 返回当前 Worker 监听的任务队列名。