切换主题
业务接入与投影
问题:创建业务单成功,平台任务还没出现
用户创建图像处理单,接口成功表示业务意图已保存,但远端调度平台暂时不可用。若直接回滚全部业务,短暂故障会让用户反复提交;若无记录地后台重试,则又会丢失请求。
让业务单拥有“提交中”状态,并持久化平台提交意图,既表达真实状态,也建立可靠恢复入口。成功语义应对用户说清,而非假装整个分布式链路已经同步完成。
原理:各系统拥有自己的事实
业务系统拥有用户、资源和业务决策;执行平台拥有编排、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;转码平台只接收不透明资源引用与流程契约,返回稳定执行标识。回调验签、去重并更新当前执行投影;查询补偿。产物发布仍由资产平台判断,不因收到成功就无条件公开。