Skip to content

业务接入与投影 ​

问题:创建业务单成功,平台任务还没出现 ​

用户创建图像处理单,接口成功表示业务意图已保存,但远端调度平台暂时不可用。若直接回滚全部业务,短暂故障会让用户反复提交;若无记录地后台重试,则又会丢失请求。

让业务单拥有“提交中”状态,并持久化平台提交意图,既表达真实状态,也建立可靠恢复入口。成功语义应对用户说清,而非假装整个分布式链路已经同步完成。

原理:各系统拥有自己的事实 ​

业务系统拥有用户、资源和业务决策;执行平台拥有编排、Attempt 与执行事实。用 stable external_ref、task_id、execution_id 关联。平台回调影响业务投影,不直接改业务数据库。

防腐层把外部 DTO 映射为内部模型,避免把远端全部状态枚举和字段散播到每个业务模块。业务可以只关心“提交中、处理中、完成、需要处理”,仍保留平台原始标识与诊断信息。

最小例子:可靠提交与双向补偿 ​

text
事务 1:业务单 SUBMITTING + command_outbox(SUBMIT, business_id)
后台提交:stable idempotency key → 执行平台
事务 2:保存 task_id / execution_id,投影为 PROCESSING
回调:验证 → Inbox + 当前执行投影(同事务)
补偿:按已知标识查询执行平台,按版本/序列合并

如果回调先于提交响应到达,可以根据稳定 external_ref 找到业务单,不能要求 task_id 必须已经写入才允许关联。

方案比较 ​

同步等待整次执行容易超过请求时间且无法可靠恢复;直接调用并立即保存远端 ID 有未知结果窗口;业务 Outbox 结合平台幂等创建可以安全重试。

回调降低查询延迟与负载,但不能假定永不遗漏;定期查询补偿提高收敛能力,代价是扫描和限流。两者都应通过同一版本规则更新投影,不能一个防倒退、另一个无条件覆盖。

真实案例:处理单与平台请求分离 ​

App Demo 创建用例检查资源已就绪与处理模式,Repository.CreateOrder 保存业务意图;提交在后台推进。TDP CreateTask 则建立另一套平台对象。

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

业务创建与命令提交边界 app-demo · c383860e

internal/core/task/application/service.go · 第 50–111 行
符号:Service.Create · 核对日期 2026-10-02
来源与提交版本一致

func (s *Service) Create(ctx context.Context, name, sourceAssetID string, processingMode taskdomain.ProcessingMode, watermark, idempotencyKey string) (*taskdomain.Order, error) {
	name, watermark = strings.TrimSpace(name), strings.TrimSpace(watermark)
	if name == "" || len(name) > 200 {
		return nil, errors.New("name must contain 1 to 200 characters")
	}
	if !processingMode.Valid() {
		return nil, errors.New("processing_mode must be ARCHIVE_CLEAN_ONLY or ARCHIVE_CLEAN_THEN_WATERMARK")
	}
	if processingMode == taskdomain.ProcessingModeArchiveCleanThenWatermark && (watermark == "" || len([]rune(watermark)) > 100) {
		return nil, errors.New("watermark_text must contain 1 to 100 characters")
	}
	if processingMode == taskdomain.ProcessingModeArchiveCleanOnly && watermark != "" {
		return nil, errors.New("watermark_text must be empty when processing_mode is ARCHIVE_CLEAN_ONLY")
	}
	asset, err := s.repository.GetAsset(ctx, sourceAssetID)
	if err != nil {
		return nil, err
	}
	if asset.Kind != assetdomain.KindSourceZIP || asset.State != "READY" {
		return nil, errors.New("source asset is not a ready ZIP")
	}
	requestHash := hashRequest(map[string]string{"name": name, "processing_mode": string(processingMode), "source_asset_id": sourceAssetID, "watermark_text": watermark})
	if replay, found, replayErr := s.replay(ctx, "CREATE", idempotencyKey, requestHash); found || replayErr != nil {
		return replay, replayErr
	}
	id, now := uuid.NewString(), time.Now().UTC()
	order, err := s.repository.CreateOrder(ctx, taskdomain.Order{
		ID:                  id,
		Name:                name,
		SourceAssetID:       sourceAssetID,
		OutputCollectionRef: assetdomain.CollectionReference(id),
		ProcessingMode:      processingMode,
		WatermarkText:       watermark,
		BusinessStatus:      taskdomain.StatusSubmitting,
		Result:              []byte(`{}`),
		CreatedAt:           now,
		UpdatedAt:           now,
	}, idempotencyKey, requestHash)
	if err == nil {
		return order, nil
	}
	if replay, found, replayErr := s.replay(ctx, "CREATE", idempotencyKey, requestHash); found || replayErr != nil {
		return replay, replayErr
	}
	return nil, err
}

func (s *Service) Get(ctx context.Context, id string) (*taskdomain.Order, error) {
	return s.repository.GetOrder(ctx, id)
}
func (s *Service) List(ctx context.Context) ([]taskdomain.Order, int64, error) {
	return s.repository.ListOrders(ctx, 100, 0)
}
func (s *Service) Cancel(ctx context.Context, id, key string) (*taskdomain.Order, error) {
	requestHash := hashRequest(map[string]string{"operation": "cancel", "order_id": id})
	if replay, found, err := s.replay(ctx, "CANCEL", key, requestHash); found || err != nil {
		return replay, err
	}
	current, err := s.repository.GetOrder(ctx, id)
	if err != nil {
		return nil, err
	}

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

创建用例与原子持久化 backend · c6dbe05c

server/internal/core/execution/application/create_task.go · 第 67–211 行
符号:CreateTask.Handle · 核对日期 2026-10-02
来源与提交版本一致

// Handle 的业务流程:校验应用、输入、回调和幂等键,解析已发布流程,
// 合并步骤调度约束并校验输入契约,最后原子写入 Task、首个 Execution、创建事件和 Outbox。
func (handler *CreateTask) Handle(ctx context.Context, command CreateTaskCommand) (*CreateTaskResult, error) {
	applicationID, err := domain.ParseID(command.ApplicationID)
	if err != nil {
		return nil, fmt.Errorf("application ID: %w", err)
	}
	input, err := domain.ParseDocument(command.Input)
	if err != nil {
		return nil, fmt.Errorf("input: %w", err)
	}
	if err := validateCallbackURL(command.CallbackURL, handler.allowInsecureCallbacks); err != nil {
		return nil, err
	}
	if strings.TrimSpace(command.PrincipalID) == "" || strings.TrimSpace(command.IdempotencyKey) == "" {
		return nil, errors.New("principal ID and idempotency key are required")
	}
	workflow, err := handler.workflows.ResolvePublished(ctx, strings.TrimSpace(command.WorkflowName), strings.TrimSpace(command.WorkflowVersion))
	if err != nil {
		return nil, fmt.Errorf("resolve workflow: %w", err)
	}
	constraints := command.StepConstraints
	if len(constraints) == 0 {
		constraints = []byte(`{}`)
	}
	workflowSnapshot, err := mergeStepConstraints(workflow.Snapshot, constraints)
	if err != nil {
		return nil, fmt.Errorf("%w: %v", fault.ErrInvalid, err)
	}
	if handler.inputValidator != nil {
		var definition struct {
			InputSchema json.RawMessage `json:"input_schema"`
		}
		if err := json.Unmarshal(workflow.Snapshot, &definition); err != nil {
			return nil, fmt.Errorf("decode workflow input contract: %w", err)
		}
		if err := handler.inputValidator.ValidateDocument(definition.InputSchema, command.Input, "workflow input"); err != nil {
			return nil, fmt.Errorf("%w: %v", fault.ErrInvalid, err)
		}
	}

	ids, err := generateIDs(handler.ids, 4)
	if err != nil {
		return nil, err
	}
	now := handler.clock.Now().UTC()
	task, execution, err := domain.NewTask(domain.NewTaskParams{
		TaskID:               ids[0],
		ExecutionID:          ids[1],
		ApplicationID:        applicationID,
		Name:                 command.Name,
		ExternalRef:          command.ExternalRef,
		WorkflowDefinitionID: workflow.ID,
		Input:                input,
		WorkflowSnapshot:     workflowSnapshot,
		StepConstraints:      constraints,
		CallbackURL:          command.CallbackURL,
		Priority:             command.Priority,
		CreatedAt:            now,
	})
	if err != nil {
		return nil, err
	}
	eventData, err := json.Marshal(map[string]any{
		"task_id":      string(task.ID),
		"execution_id": string(execution.ID),
		"external_ref": task.ExternalRef,
		"status":       execution.Status,
	})
	if err != nil {
		return nil, fmt.Errorf("encode execution event: %w", err)
	}
	event := domain.Event{
		ID:            ids[2],
		OutboxID:      ids[3],
		ApplicationID: applicationID,
		ExecutionID:   execution.ID,
		Sequence:      1,
		Type:          domain.EventExecutionCreated,
		Data:          eventData,
		OccurredAt:    now,
	}
	requestHash := sha256.Sum256([]byte(strings.Join([]string{
		command.ApplicationID, command.Name, command.ExternalRef, command.WorkflowName, command.WorkflowVersion, string(command.Input), string(constraints), command.CallbackURL, fmt.Sprint(command.Priority),
	}, "\x00")))
	persistedTask, persistedExecution, err := handler.repository.Create(ctx, task, execution, event, ports.Idempotency{
		PrincipalID: command.PrincipalID,
		Operation:   "create-task",
		Key:         command.IdempotencyKey,
		RequestHash: requestHash[:],
	})
	if err != nil {
		return nil, fmt.Errorf("persist task: %w", err)
	}
	return &CreateTaskResult{
		Task:      persistedTask,
		Execution: persistedExecution,
	}, nil
}

func mergeStepConstraints(snapshot, raw []byte) ([]byte, error) {
	var document map[string]any
	var overrides map[string]map[string]any
	if err := json.Unmarshal(snapshot, &document); err != nil {
		return nil, errors.New("invalid workflow snapshot")
	}
	if err := json.Unmarshal(raw, &overrides); err != nil {
		return nil, errors.New("invalid step constraints")
	}
	steps, ok := document["steps"].([]any)
	if !ok {
		return nil, errors.New("workflow steps are missing")
	}
	known := make(map[string]bool, len(steps))
	for _, value := range steps {
		step, ok := value.(map[string]any)
		if !ok {
			return nil, errors.New("invalid workflow step")
		}
		key, _ := step["key"].(string)
		known[key] = true
		override, exists := overrides[key]
		if !exists {
			continue
		}
		affinity, _ := step["affinity"].(map[string]any)
		if affinity == nil {
			affinity = map[string]any{}
		}
		for _, field := range []string{"node_id", "group"} {
			if value, exists := override[field]; exists {
				affinity[field] = value
			}
		}
		for _, field := range []string{"required_labels", "preferred_labels", "minimum_resources"} {
			if value, exists := override[field]; exists {
				base, _ := affinity[field].(map[string]any)
				if base == nil {
					base = map[string]any{}
				}
				values, ok := value.(map[string]any)
				if !ok {
					return nil, fmt.Errorf("step %s %s must be an object", key, field)
				}
				for k, v := range values {

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

Callback 校验、去重与缺口 app-demo · c383860e

internal/core/task/application/callback.go · 第 26–109 行
符号:ReceiveCallback / VerifyCallbackSignature · 核对日期 2026-10-02
来源与提交版本一致

func (s *Service) ReceiveCallback(ctx context.Context, input taskdomain.CallbackInput, applicationID string, discard bool, now time.Time) error {
	if input.SpecVersion != "1.0" || input.Sequence < 1 || input.EventID == "" || input.Type == "" {
		return ErrInvalidCallback
	}
	if parsed, err := uuid.Parse(input.EventID); err != nil || parsed == uuid.Nil {
		return ErrInvalidCallback
	}
	if input.ApplicationID != "" && input.ApplicationID != applicationID {
		return ErrInvalidCallback
	}
	if parsed, err := uuid.Parse(input.ExecutionID); err != nil || parsed == uuid.Nil {
		return ErrInvalidCallback
	}
	if input.TaskID != "" {
		if parsed, err := uuid.Parse(input.TaskID); err != nil || parsed == uuid.Nil {
			return ErrInvalidCallback
		}
	}
	if input.ExternalRef == "" {
		order, err := s.repository.GetOrderByExecution(ctx, input.ExecutionID)
		if err != nil {
			return ErrInvalidCallback
		}
		input.ExternalRef = order.ID
	}
	if parsed, err := uuid.Parse(input.ExternalRef); err != nil || parsed == uuid.Nil {
		return ErrInvalidCallback
	}
	if input.Status != "" && !map[string]bool{"CREATED": true, "QUEUED": true, "RUNNING": true, "SUCCEEDED": true, "FAILED": true, "TIMED_OUT": true, "CANCELING": true, "CANCELED": true}[input.Status] {
		return ErrInvalidCallback
	}
	if discard {
		s.logger.Warn("demo discarded callback after validation", "event_id", input.EventID, "execution_id", input.ExecutionID, "sequence", input.Sequence)
		return nil
	}
	progress := 0
	if input.Progress != nil {
		progress = *input.Progress
	}
	_, gap, err := s.repository.AcceptCallback(ctx, taskports.CallbackEvent{
		ID:          input.EventID,
		TaskID:      input.TaskID,
		ExecutionID: input.ExecutionID,
		Type:        input.Type,
		Sequence:    input.Sequence,
		Body:        append([]byte(nil), input.RawBody...),
		ReceivedAt:  now,
	}, input.ExternalRef, input.Status, int32(progress), input.Data, now)
	if gap {
		s.logger.Warn("callback sequence gap", "order_id", input.ExternalRef, "execution_id", input.ExecutionID, "sequence", input.Sequence)
	}
	return err
}

func VerifyCallbackSignature(header string, body, secret []byte, now time.Time) error {
	timestampPart, signaturePart, ok := strings.Cut(header, ",")
	if !ok {
		return ErrInvalidCallbackSignature
	}
	rawTimestamp, okTimestamp := strings.CutPrefix(timestampPart, "t=")
	rawSignature, okSignature := strings.CutPrefix(signaturePart, "v1=")
	if !okTimestamp || !okSignature {
		return ErrInvalidCallbackSignature
	}
	timestamp, err := strconv.ParseInt(rawTimestamp, 10, 64)
	if err != nil {
		return ErrInvalidCallbackSignature
	}
	if delta := now.Sub(time.Unix(timestamp, 0)); delta < -5*time.Minute || delta > 5*time.Minute {
		return ErrExpiredCallback
	}
	received, err := hex.DecodeString(rawSignature)
	if err != nil {
		return ErrInvalidCallbackSignature
	}
	mac := hmac.New(sha256.New, secret)
	_, _ = fmt.Fprintf(mac, "%d.", timestamp)
	_, _ = mac.Write(body)
	if !hmac.Equal(received, mac.Sum(nil)) {
		return ErrInvalidCallbackSignature
	}
	return nil
}

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

跨项目端到端资料 workbench · 514f5a00

docs/learning/README.md · 第 11–57 行
符号:主链路与事实边界 · 核对日期 2026-10-02
来源与提交版本一致

## 最有效的学习方式

不要从目录树第一行开始顺序读完整个项目。先跟一条真实的 `ARCHIVE_CLEAN_THEN_WATERMARK` 处理单纵向走完,再回头学习每个模块的横向设计。

这条纵向链路覆盖了项目最重要的机制:

```text
浏览器上传 ZIP
  → app-demo 创建 ProcessingOrder 和本地 Command Outbox
  → app-demo 后台提交 TDP Task
  → TDP 创建 Task / Execution / Event / Outbox
  → NATS JetStream 唤醒 Temporal
  → Temporal 展开 DAG 并创建 Step / Attempt
  → PostgreSQL Command 经 SSE 到 Worker
  → Worker SQLite Inbox → archive.clean → 本地 Event Outbox
  → Server Attempt Inbox → Temporal Signal
  → archive.clean 输出经 CEL 映射给 image.watermark
  → Execution 完成并可靠 Callback
  → app-demo Callback Inbox 更新业务投影
```

完整细节见[一条处理单的数据流](business-flow.md)。

## 先建立四个“事实边界”

| 边界 | 保存内容 | 是否可由其他系统直接改写 |
|---|---|---|
| app-demo PostgreSQL | 处理单、业务状态、资源元数据、提交 Outbox、Callback Inbox | 否;TDP 只能通过 API/Callback 影响业务投影 |
| TDP PostgreSQL | Task、Execution、Step、Attempt、Worker、租约、事件和可靠消息 | 否;Temporal、NATS 和 Worker 都不能绕过应用用例直接成为事实源 |
| Worker SQLite | 当前 Worker 的命令 Inbox、Attempt、本地 Checkpoint、事件 Outbox、注册身份 | 否;只属于该 Worker,不是全局查询库 |
| MinIO | app-demo 拥有的 ZIP、PNG 和 manifest 文件字节 | 通过短期 GET/PUT 授权访问;TDP Server 不转发文件内容 |

Temporal 保存耐久编排历史,NATS 保存可靠传输中的消息;它们都不替代 PostgreSQL 查询模型。删除 Worker SQLite 也不是“清缓存”,而是放弃该实例尚未同步的恢复状态。

## 文档地图

1. [一条处理单的数据流](business-flow.md)
   - 以 `ARCHIVE_CLEAN_THEN_WATERMARK` 为主线,从上传、创建、调度、Worker 执行一直跟到 Callback。
   - 重点看数据在业务 ID、资源引用、AccessGrant、Step 输出和业务投影之间如何变形。
2. [数据表与一致性边界](data-model.md)
   - 分别解释 app-demo PostgreSQL、TDP PostgreSQL、Worker SQLite 的表设计。
   - 不只列字段,还说明每组表为何要在同一个事务里写、哪些列是事实、哪些列是投影。
3. [代码阅读与联调练习](code-reading.md)
   - 给出按业务动作定位代码的方法,以及一组不会删除本地数据的观察命令。
   - 用同一组 ID 串起页面、日志、两套 PostgreSQL、Temporal UI 和 NATS NUI。

已有文档仍然是重要的专题参考:

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

任务链路把 ZIP 清理的输出映射给水印步骤。业务资源引用稳定,短期授权按需要签发,最终 Callback 更新处理单投影。平台不保存用户业务文件字节,也不拥有水印文字的业务意义。

失败与边界:重跑后的迟到回调 ​

同一业务单可能有多次 Execution。旧运行的成功消息晚到不能覆盖当前新运行状态。关联当前 execution_id、校验序列并保留历史,不把 task_id 当作唯一运行身份。

Callback 失败要可重投,投影更新失败不应成功 ACK;确实重复且已提交的事件可以成功回应。长期无法收敛时,需要显示可诊断中间状态,不能永远只有旋转加载图标。

迁移练习与参考答案 ​

练习:资产平台接入转码服务,转码服务不认识租户业务规则。如何划边界?

参考答案:资产平台检查用户与资源权限,保存业务单和提交 Outbox;转码平台只接收不透明资源引用与流程契约,返回稳定执行标识。回调验签、去重并更新当前执行投影;查询补偿。产物发布仍由资产平台判断,不因收到成功就无条件公开。

继续阅读:端到端案例、Inbox 与幂等。

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