Skip to content

任务模型与状态机 ​

问题:失败后把状态改回待执行就够了吗 ​

报表导出失败后,用户点击重新执行。如果只把同一行从 FAILED 改成 PENDING,上次用了哪版模板、在哪一步失败、是否已生成部分文件都会被覆盖。迟到的旧结果还可能被当成新运行成功。

长期意图和一次运行应有不同身份。进一步把运行拆成步骤,把每次实际执行拆成尝试,可以分别回答“用户要什么”“本次怎样做”“哪个动作”“第几次做”。

原理:状态是对事实的压缩表达 ​

状态机明确哪些转换被允许,转换需要什么证据,以及终态是否可以再被普通事件改变。不要用一组可随意组合的 running/failed/canceled 布尔值,让矛盾状态变得可表示。

取消应区分“已提出取消”与“确认停止”。已经开始的外部动作无法总是立即撤销。CANCELING 是中间事实,CANCELED 需要确认或明确的失效判定。

最小例子:状态转换是一个领域行为 ​

go
// 教学片段:持久化版本检查在 repository 中完成。
func (run *Run) Finish(result Result) error {
    if run.State != "RUNNING" { return errors.New("run is not running") }
    run.State = "SUCCEEDED"
    run.Result = result
    run.Version++
    return nil
}

状态方法保护单对象规则;保存时必须验证读到的版本仍有效。同一次成功事件重投应返回同一结果;失败或取消后收到不同终态要遵循明确竞争策略。

交互 04

取消是请求,终态是事实

Execution #1 · CREATED
  1. Execution #1 · CREATED

转换集合来自 Execution 领域模型。重跑入口在本教学模型中限制为终态,具体权限与业务条件仍由应用用例决定。这里没有模拟 Worker 停止、定时器或外部副作用;CANCELED 的落库必须有相应证据。

尝试先进入 CANCELING,观察为什么不能再直接 SUCCEEDED;再创建新执行,观察上次状态历史如何保留。模型允许某些早期状态直接终结,这是当前案例的转换规则,其他业务可以不同。

方案比较 ​

单一 job 行适合简单、短生命周期工作,但重试审计和步骤恢复困难。Task/Run/Step/Attempt 分层提供历史与幂等边界,也增加查询和状态聚合成本。不要为了统一命名强迫每个通知任务都建立四层。

状态机可由代码枚举、数据驱动定义或耐久工作流引擎表达。数据库约束、领域方法与编排器必须保持一致;图中有一条箭头不代表任意时刻都可以通过 API 强制转换。

真实案例:四层执行模型 ​

TDP 的 Task 保存长期请求及 CurrentExecutionID;Execution 保存运行序号和 workflow_snapshot;Step 表达 DAG 节点;Attempt 对应一次 Worker 执行。内部重试创建新 Attempt,用户重跑产生新的 Execution。

已核对的实现 · 本地代码快照

实体、值对象与创建不变量 backend · c6dbe05c

server/internal/core/execution/domain/model.go · 第 12–141 行
符号:NewTask · 核对日期 2026-10-02
来源与提交版本一致

// ID 是 Execution 上下文使用的强类型资源标识,避免把普通字符串误传为领域 ID。
type ID string

// ParseID 去除首尾空白并拒绝空标识。
func ParseID(value string) (ID, error) {
	value = strings.TrimSpace(value)
	if value == "" {
		return "", errors.New("ID is required")
	}
	return ID(value), nil
}

// Document 保存已经校验过的 JSON 文档,并在进入领域层时复制以避免外部修改。
type Document []byte

// ParseDocument 校验并复制 JSON 字节,调用方可安全复用原缓冲区。
func ParseDocument(value []byte) (Document, error) {
	if len(value) == 0 || !json.Valid(value) {
		return nil, errors.New("document must contain valid JSON")
	}
	return append(Document(nil), value...), nil
}

// Task 是业务请求的长期聚合根;一次 Task 可以通过重跑产生多个 Execution。
type Task struct {
	ID                   ID              // 平台生成的任务标识。
	ApplicationID        ID              // 任务的归属和访问范围。
	Name                 string          // 供业务人员识别的任务名称。
	ExternalRef          string          // 业务系统用于关联和幂等追踪的外部引用。
	WorkflowDefinitionID ID              // 创建任务时解析到的已发布流程。
	Input                Document        // 业务系统提交给流程的原始输入。
	StepConstraints      Document        // 按步骤合并后的调度约束。
	EffectiveStatus      ExecutionStatus // 当前执行状态,便于直接查询任务。
	CurrentExecutionID   ID              // 最近一次创建或重跑的执行。
	Priority             int16           // 待调度任务的业务优先级,范围为 -100 到 100。
	Version              int64           // 用于任务的乐观并发控制。
	CreatedAt            time.Time       // 任务首次创建时间。
	UpdatedAt            time.Time       // 聚合最后一次状态变更时间。
}

// Execution 是 Task 的一次具体运行,保存创建时冻结的流程快照。
type Execution struct {
	ReconciliationError string          // 终态对账失败或待人工检查的原因。
	SchedulingReasons   []string        // 当前排队步骤的等待原因。
	ID                  ID              // 本次执行的稳定标识。
	TaskID              ID              // 所属任务。
	Number              int32           // 同一任务下从 1 开始递增的执行序号。
	Status              ExecutionStatus // 执行生命周期状态。
	Progress            int16           // 0 到 100 的单调进度值。
	Result              Document        // 执行完成后的 JSON 结果。
	WorkflowSnapshot    Document        // 冻结本次运行使用的已编译流程。
	CallbackURL         string          // 完成事件的业务回调地址。
	Reason              string          // 记录重跑、取消或失败等人工可读原因。
	Version             int64           // 用于执行的乐观并发控制。
	CreatedAt           time.Time       // 执行创建时间。
	StartedAt           *time.Time      // 首次进入运行态的时间,未开始时为空。
	CompletedAt         *time.Time      // 进入终态的时间,未结束时为空。
}

// NewTaskParams 汇集原子创建 Task 及首个 Execution 所需的已解析参数。
type NewTaskParams struct {
	TaskID               ID        // 预先生成的任务标识。
	ExecutionID          ID        // 首个执行标识。
	ApplicationID        ID        // 任务所属应用。
	Name                 string    // 任务展示名称。
	ExternalRef          string    // 业务系统外部引用。
	WorkflowDefinitionID ID        // 已发布流程标识。
	Input                Document  // 已经校验的任务输入。
	WorkflowSnapshot     Document  // 创建时冻结的流程定义。
	StepConstraints      Document  // 已经合并的步骤调度约束。
	CallbackURL          string    // 可选的执行事件回调地址。
	Priority             int16     // 任务优先级。
	CreatedAt            time.Time // 统一注入的创建时间。
}

// NewTask 校验创建不变量,并同时构造 Task 与编号为 1 的 Execution。
func NewTask(params NewTaskParams) (*Task, *Execution, error) {
	if err := validateIDs(params.TaskID, params.ExecutionID, params.ApplicationID, params.WorkflowDefinitionID); err != nil {
		return nil, nil, err
	}
	if strings.TrimSpace(params.Name) == "" {
		return nil, nil, errors.New("task name is required")
	}
	if len([]rune(strings.TrimSpace(params.Name))) > 200 {
		return nil, nil, errors.New("task name must not exceed 200 characters")
	}
	if strings.TrimSpace(params.ExternalRef) == "" {
		return nil, nil, errors.New("external reference is required")
	}
	if params.Priority < -100 || params.Priority > 100 {
		return nil, nil, errors.New("priority must be between -100 and 100")
	}
	if len(params.StepConstraints) == 0 {
		params.StepConstraints = Document(`{}`)
	}
	if !json.Valid(params.Input) || !json.Valid(params.WorkflowSnapshot) || !json.Valid(params.StepConstraints) {
		return nil, nil, errors.New("input and workflow snapshot must contain valid JSON")
	}
	if params.CreatedAt.IsZero() {
		return nil, nil, errors.New("creation time is required")
	}

	execution := &Execution{
		ID:               params.ExecutionID,
		TaskID:           params.TaskID,
		Number:           1,
		Status:           ExecutionCreated,
		WorkflowSnapshot: append(Document(nil), params.WorkflowSnapshot...),
		CallbackURL:      params.CallbackURL,
		Version:          1,
		CreatedAt:        params.CreatedAt,
	}
	task := &Task{
		ID:                   params.TaskID,
		ApplicationID:        params.ApplicationID,
		Name:                 strings.TrimSpace(params.Name),
		ExternalRef:          strings.TrimSpace(params.ExternalRef),
		WorkflowDefinitionID: params.WorkflowDefinitionID,
		Input:                append(Document(nil), params.Input...),
		StepConstraints:      append(Document(nil), params.StepConstraints...),
		EffectiveStatus:      execution.Status,
		CurrentExecutionID:   execution.ID,
		Priority:             params.Priority,
		Version:              1,
		CreatedAt:            params.CreatedAt,
		UpdatedAt:            params.CreatedAt,
	}
	return task, execution, nil
}

片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。

执行状态机及单调进度 backend · c6dbe05c

server/internal/core/execution/domain/status.go · 第 22–79 行
符号:Transition / SetProgress · 核对日期 2026-10-02
来源与提交版本一致

var executionTransitions = map[ExecutionStatus]map[ExecutionStatus]struct{}{
	ExecutionCreated:   set(ExecutionQueued, ExecutionSucceeded, ExecutionFailed, ExecutionTimedOut, ExecutionCanceling, ExecutionCanceled),
	ExecutionQueued:    set(ExecutionRunning, ExecutionSucceeded, ExecutionFailed, ExecutionTimedOut, ExecutionCanceling, ExecutionCanceled),
	ExecutionRunning:   set(ExecutionSucceeded, ExecutionFailed, ExecutionTimedOut, ExecutionCanceling, ExecutionCanceled),
	ExecutionCanceling: set(ExecutionCanceled, ExecutionFailed, ExecutionTimedOut),
}

// Terminal 判断执行是否已经结束且不会再接受普通状态推进。
func (s ExecutionStatus) Terminal() bool {
	switch s {
	case ExecutionSucceeded, ExecutionFailed, ExecutionTimedOut, ExecutionCanceled:
		return true
	default:
		return false
	}
}

// Transition 按状态机推进执行,并统一维护开始时间、完成时间和版本号。
func (e *Execution) Transition(to ExecutionStatus, at time.Time) error {
	if e.Status == to {
		return nil
	}
	if _, ok := executionTransitions[e.Status][to]; !ok {
		return fmt.Errorf("execution transition %s -> %s is not allowed", e.Status, to)
	}
	if to == ExecutionRunning && e.StartedAt == nil {
		startedAt := at
		e.StartedAt = &startedAt
	}
	if to.Terminal() {
		completedAt := at
		e.CompletedAt = &completedAt
		if to == ExecutionSucceeded {
			e.Progress = 100
		}
	}
	e.Status = to
	e.Version++
	return nil
}

// SetProgress 仅允许运行中的执行单调增加 0 到 100 的进度。
func (e *Execution) SetProgress(progress int16) error {
	if e.Status != ExecutionRunning {
		return fmt.Errorf("cannot update progress while execution is %s", e.Status)
	}
	if progress < e.Progress || progress > 100 {
		return fmt.Errorf("progress must be monotonic and between %d and 100", e.Progress)
	}
	e.Progress = progress
	e.Version++
	return nil
}

// StepStatus 表示工作流步骤从等待依赖到完成的生命周期状态。
type StepStatus string

const (

片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。

Execution.Transition 对同状态重复更新直接返回,不重复递增版本;进入 RUNNING 设置首次开始时间,进入终态设置完成时间。SetProgress 只在 RUNNING 允许单调增加,成功终态设置为 100。

失败与边界:进度不等于业务成功 ​

进度可以是估算值,100% 的文件处理仍可能在保存结果时失败。成功应以持久化终态及必要输出为依据。多个并发步骤的总体进度需明确权重,避免重试导致进度倒退或重复累计。

超时说明协调者不再等待,并不证明远端进程已结束。Attempt 租约、事件来源版本和资源回收需要共同管理。终态不可回退,也不意味着不再做对账或清理。

迁移练习与参考答案 ​

练习:视频转码有下载、编码、上传三步;编码失败重试,用户又发起整次重跑。如何分配标识?

参考答案:视频业务请求身份稳定;每次整次运行一个 Run;三步分别有 Step;编码每次重试有不同 Attempt。产物路径或远端请求应含相应执行身份,旧 Attempt 的迟到事件不能覆盖新 Attempt。保留旧 Run 失败记录供诊断。

继续阅读:DAG 与耐久编排、调度与资源协调。

理解原理 · 分析取舍 · 用真实代码检验