Skip to content

一条任务的旅程 ​

这个案例用于检验课程中的知识。业务是上传 ZIP,清理不需要的文件,再给图片加水印。关注每段链路的所有权、原子边界和恢复证据,而不是记忆项目路径。

总体链路 ​

从业务意图到执行结果
100%
业务处理单
Command Outbox
→
TDP Task / Execution
Event / Outbox
NATS → Temporal DAG
→
Step / Attempt
Worker Command
Worker SQLite Inbox
Capability / Event Outbox
→
Server Attempt Inbox
状态推进 / Signal
执行事件 → Callback
→
业务 Callback Inbox
处理单投影

可触控平移或键盘滚动;每张图的正文同时提供文字解释。

1. 上传与创建:先保存业务意图 ​

浏览器上传资源,业务系统保存资源元数据,字节进入对象存储。创建处理单要求源 ZIP 已 READY;水印模式要求合法文本,纯清理模式拒绝多余水印参数。

业务创建成功后可能仍处于 SUBMITTING。这不是系统错误,而是“本地意图已保存、平台提交还在推进”的真实中间状态。处理单和 Command Outbox 建立可恢复提交依据。

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

业务创建与命令提交边界 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
	}

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

知识迁移:订单与支付请求、文章与索引请求、资产与转码请求都可有这种本地意图与外部推进的分离。

2. 平台创建:冻结本次执行 ​

后台用稳定幂等键调用 TDP。CreateTask 解析发布 Workflow、输入和步骤约束,构造 Task、首个 Execution、事件与请求指纹,然后通过 Repository 原子保存。

Task 是长期请求,Execution 保存本次运行的流程快照。外部提交超时可以重复调用,不能更换键创建另一个任务。

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

创建用例与原子持久化 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 {

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

核心定义协作契约 backend · c6dbe05c

server/internal/core/execution/ports/outbound/repository.go · 第 9–24 行
符号:Repository · 核对日期 2026-10-02
来源与提交版本一致

// Idempotency 描述一次写操作的调用方、键和请求指纹。
type Idempotency struct {
	PrincipalID string // 隔离不同调用方的同名幂等键。
	Operation   string // 隔离同一调用方的不同业务动作。
	Key         string // 调用方提供的幂等键。
	RequestHash []byte // 用于拒绝复用同一键但内容不同的请求。
}

// Repository 原子持久化 Task、Execution、首个事件和幂等记录。
type Repository interface {
	// Create 创建聚合;重复且请求指纹一致时返回原结果。
	Create(context.Context, *domain.Task, *domain.Execution, domain.Event, Idempotency) (*domain.Task, *domain.Execution, error)
	// Get 在应用范围内读取任务。
	Get(context.Context, Scope, domain.ID) (*domain.Task, error)
}

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

故障窗口:平台已保存但业务未收到响应。恢复依赖平台幂等结果;回调提前到达也可用 external_ref 关联业务单。

3. Outbox 发布:数据库事实跨越消息边界 ​

平台 Outbox 在短事务内认领,发布执行事件到 NATS JetStream,再确认发布。接收侧桥接到 Temporal,稳定标识吸收重复推进。

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

发布、认领与退避 backend · c6dbe05c

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

const (
	outboxClaimTTL    = time.Minute
	outboxMaxAttempts = 20
)

// OutboxDispatcher 将领域事务中的 Outbox 事件可靠发布到消息系统。
type OutboxDispatcher struct {
	store     executionports.OutboxStore
	publisher executionports.EventPublisher
	clock     executionports.Clock
}

func NewOutboxDispatcher(store executionports.OutboxStore, publisher executionports.EventPublisher, clock executionports.Clock) (*OutboxDispatcher, error) {
	if store == nil || publisher == nil || clock == nil {
		return nil, errors.New("outbox store, event publisher, and clock are required")
	}
	return &OutboxDispatcher{
		store:     store,
		publisher: publisher,
		clock:     clock,
	}, nil
}

func (dispatcher *OutboxDispatcher) DispatchOutboxBatch(ctx context.Context, limit int) (int, []error) {
	if limit < 1 {
		return 0, []error{errors.New("outbox batch limit must be positive")}
	}
	now := dispatcher.clock.Now().UTC()
	events, err := dispatcher.store.ClaimOutboxEvents(ctx, now, outboxClaimTTL, limit)
	if err != nil {
		return 0, []error{fmt.Errorf("claim outbox events: %w", err)}
	}
	failures := make([]error, 0)
	for _, event := range events {
		err := dispatcher.publisher.PublishExecutionEvent(ctx, executionports.PublishedEvent{
			ID:            event.ID,
			Type:          event.EventType,
			CorrelationID: event.AggregateID,
			Sequence:      event.Sequence,
			OccurredAt:    event.CreatedAt,
			Payload:       event.Payload,
		})
		if err == nil {
			err = dispatcher.store.MarkOutboxEventPublished(ctx, event.ID, event.ClaimID, dispatcher.clock.Now().UTC())
			if err != nil {
				failures = append(failures, fmt.Errorf("mark outbox event %s published: %w", event.ID, err))
			}
			continue
		}
		status := "FAILED"
		if event.Attempts >= outboxMaxAttempts {
			status = "DEAD"
		}
		availableAt := dispatcher.clock.Now().UTC().Add(time.Second * time.Duration(1<<min(event.Attempts, 10)))
		if markErr := dispatcher.store.MarkOutboxEventFailed(ctx, event.ID, event.ClaimID, status, availableAt, err.Error()); markErr != nil {
			err = errors.Join(err, markErr)
		}
		failures = append(failures, fmt.Errorf("publish outbox event %s: %w", event.ID, err))
	}
	return len(events), failures
}

func (dispatcher *OutboxDispatcher) NextOutboxAttemptAt(ctx context.Context) (time.Time, bool, error) {
	next, ok, err := dispatcher.store.NextOutboxAttemptAt(ctx)
	if err != nil {
		return time.Time{}, false, fmt.Errorf("read next outbox attempt: %w", err)
	}
	return next, ok, nil
}

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

Outbox 认领与所有权确认 backend · c6dbe05c

server/db/queries/runtime.sql · 第 4–50 行
符号:ClaimOutboxEvents / MarkOutboxEventPublished · 核对日期 2026-10-02
来源与提交版本一致

-- name: ClaimOutboxEvents :many
WITH candidates AS (
    SELECT id
    FROM outbox_events
    WHERE (outbox_events.status IN ('PENDING', 'FAILED') AND outbox_events.available_at <= sqlc.arg(now))
       OR (outbox_events.status = 'PUBLISHING' AND outbox_events.claim_expires_at <= sqlc.arg(now))
    ORDER BY outbox_events.available_at, outbox_events.created_at
    FOR UPDATE SKIP LOCKED
    LIMIT sqlc.arg(batch_limit)
)
UPDATE outbox_events AS event
SET status = 'PUBLISHING',
    attempts = event.attempts + 1,
    claim_id = sqlc.arg(claim_id),
    claim_expires_at = sqlc.arg(claim_expires_at)
FROM candidates
WHERE event.id = candidates.id
RETURNING event.id, event.aggregate_id, event.event_type, event.sequence, event.payload,
          event.created_at, event.attempts, event.claim_id;

-- name: GetNextOutboxAttemptAt :one
SELECT MIN(candidate.due_at)::timestamptz AS due_at
FROM (
    (SELECT available_at AS due_at
     FROM outbox_events
     WHERE status IN ('PENDING', 'FAILED')
     ORDER BY available_at
     LIMIT 1)
    UNION ALL
    (SELECT claim_expires_at AS due_at
     FROM outbox_events
     WHERE status = 'PUBLISHING' AND claim_expires_at IS NOT NULL
     ORDER BY claim_expires_at
     LIMIT 1)
) AS candidate;

-- name: MarkOutboxEventPublished :execrows
UPDATE outbox_events
SET status = 'PUBLISHED', published_at = sqlc.arg(published_at), last_error = NULL,
    claim_id = NULL, claim_expires_at = NULL
WHERE id = sqlc.arg(id) AND status = 'PUBLISHING' AND claim_id = sqlc.arg(claim_id);

-- name: MarkOutboxEventFailed :execrows
UPDATE outbox_events
SET status = sqlc.arg(status), available_at = sqlc.arg(available_at),
    last_error = sqlc.arg(last_error), claim_id = NULL, claim_expires_at = NULL
WHERE id = sqlc.arg(id) AND status = 'PUBLISHING' AND claim_id = sqlc.arg(claim_id);

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

消息不是平台数据库事务的一部分;发送后确认前崩溃会重投。NATS 的去重窗口也不能代替全链路永久业务幂等。

4. DAG 与调度:先依赖,后资源 ​

Temporal 根据冻结定义展开清理与水印步骤。清理完成后,输出经输入映射成为水印输入。依赖和输入契约分别校验。

Step 进入可调度状态后,选择器过滤 Worker 能力、版本、标签、在线与静态资源,再打分;派发事务复核执行状态和容量,创建 Attempt 与 worker_commands。Worker 唤醒失败不丢命令,后续仍可读取持久记录。

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

DAG 依赖与环检测 backend · c6dbe05c

server/internal/adapters/inbound/temporal/dag_workflow.go · 第 45–111 行
符号:validateStepPlans · 核对日期 2026-10-02
来源与提交版本一致

func validateStepPlans(steps []StepPlan) error {
	known := make(map[string]StepPlan, len(steps))
	for _, step := range steps {
		if step.Key == "" || step.StepID == "" {
			return errors.New("step key and ID are required")
		}
		if _, exists := known[step.Key]; exists {
			return fmt.Errorf("duplicate step key %q", step.Key)
		}
		known[step.Key] = step
	}
	for _, step := range steps {
		for _, dependency := range step.DependsOn {
			if _, exists := known[dependency]; !exists {
				return fmt.Errorf("step %q depends on unknown step %q", step.Key, dependency)
			}
		}
	}
	state := make(map[string]uint8, len(steps))
	var visit func(string) error
	visit = func(key string) error {
		switch state[key] {
		case 1:
			return fmt.Errorf("dependency cycle includes step %q", key)
		case 2:
			return nil
		}
		state[key] = 1
		for _, dependency := range known[key].DependsOn {
			if err := visit(dependency); err != nil {
				return err
			}
		}
		state[key] = 2
		return nil
	}
	for key := range known {
		if err := visit(key); err != nil {
			return err
		}
	}
	return nil
}

func dependenciesComplete(step StepPlan, states []dagStepState, byKey map[string]int) bool {
	for _, dependency := range step.DependsOn {
		status := states[byKey[dependency]].status
		if status != dagStepSucceeded && status != dagStepSkipped {
			return false
		}
	}
	return true
}

func hasRunningSteps(states []dagStepState) bool {
	for _, state := range states {
		if state.status == dagStepRunning {
			return true
		}
	}
	return false
}

func firstRunningStepID(steps []StepPlan, states []dagStepState) string {
	for index, state := range states {
		if state.status == dagStepRunning {
			return steps[index].StepID

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

硬约束筛选与确定性评分 backend · c6dbe05c

server/internal/core/scheduling/domain/scheduler.go · 第 58–112 行
符号:Select · 核对日期 2026-10-02
来源与提交版本一致

// Select 先应用硬约束,再按得分降序、Node ID 升序确定唯一候选。
func Select(request Request, nodes []Node) (*Candidate, error) {
	if request.Capability == "" || request.CapabilityVersion == "" {
		return nil, errors.New("capability and version are required")
	}
	if request.Now.IsZero() || request.QueuedAt.IsZero() || request.HeartbeatTTL <= 0 {
		return nil, errors.New("queue time, scheduling time, and heartbeat TTL are required")
	}
	candidates := make([]Candidate, 0, len(nodes))
	for _, node := range nodes {
		if request.NodeID != "" && node.ID != request.NodeID {
			continue
		}
		if node.ManagementStatus != "ACTIVE" || node.Worker.RuntimeState != "READY" || node.Worker.SchedulingStatus != "ELIGIBLE" {
			continue
		}
		if node.Worker.LastHeartbeat.Before(request.Now.Add(-request.HeartbeatTTL)) || !slices.Contains(node.Worker.CapabilityVersions, request.CapabilityVersion) {
			continue
		}
		if len(QualificationReasons(node.TrustedLabels, node.Resources, request.RequiredLabels, request.MinimumResources)) > 0 {
			continue
		}
		preferred := matchCount(node.TrustedLabels, request.PreferredLabels)
		weight := node.Weight
		if weight <= 0 {
			weight = 100
		}
		score := int64(preferred*1_000_000 + weight*1_000 - node.RunningAttempts*10_000)
		candidates = append(candidates, Candidate{
			Node:     node,
			WorkerID: node.Worker.ID,
			Score:    score,
		})
	}
	if len(candidates) == 0 {
		return nil, errors.New("no eligible node")
	}
	slices.SortFunc(candidates, func(left, right Candidate) int {
		if byScore := cmp.Compare(right.Score, left.Score); byScore != 0 {
			return byScore
		}
		return cmp.Compare(left.Node.ID, right.Node.ID)
	})
	return &candidates[0], nil
}

// QualificationReasons 统一提供准入与只读诊断使用的资格过滤原因。
func QualificationReasons(labels map[string]string, capacity Resources, required map[string]string, minimum Resources) []string {
	var reasons []string
	if !matches(labels, required) {
		reasons = append(reasons, "必选标签不匹配")
	}
	if capacity.CPUMillis < minimum.CPUMillis {
		reasons = append(reasons, "CPU 静态规格不足")
	}

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

调度事务与幂等 Attempt backend · c6dbe05c

server/internal/adapters/outbound/postgres/attempt_dispatch.go · 第 45–133 行
符号:dispatchAttempt · 核对日期 2026-10-02
来源与提交版本一致

func (dispatcher *AttemptDispatcher) dispatchAttempt(ctx context.Context, request executionports.CreateAttemptRequest, claim *pgqueries.SchedulingRequest) (*executionports.AttemptRef, error) {
	stepID, err := uuid.Parse(request.StepID)
	if err != nil {
		return nil, fmt.Errorf("parse step ID: %w", err)
	}
	tx, err := dispatcher.database.BeginTx(ctx, pgx.TxOptions{
		IsoLevel: pgx.ReadCommitted,
	})
	if err != nil {
		return nil, err
	}
	defer func() { _ = tx.Rollback(ctx) }()
	txQueries := dispatcher.queries.WithTx(tx)
	lockKey := request.StepID + ":" + fmt.Sprint(request.AttemptNumber)
	if err := txQueries.AcquireTransactionAdvisoryLock(ctx, lockKey); err != nil {
		return nil, err
	}
	if _, err := txQueries.LockSchedulingRequest(ctx, pgqueries.LockSchedulingRequestParams{
		ID:      claim.ID,
		ClaimID: claim.ClaimID,
	}); err != nil {
		return nil, err
	}
	executionIDValue, err := uuid.Parse(request.ExecutionID)
	if err != nil {
		return nil, fmt.Errorf("%w: execution ID", executiondomain.ErrDispatchInvalid)
	}
	execution, err := txQueries.LockExecutionByID(ctx, executionIDValue)
	if err != nil {
		return nil, err
	}
	if execution.Status != "CREATED" && execution.Status != "QUEUED" && execution.Status != "RUNNING" {
		return nil, fmt.Errorf("%w: execution closed", executiondomain.ErrDispatchInvalid)
	}
	existing, err := txQueries.FindAttemptByStepAndNumber(ctx, pgqueries.FindAttemptByStepAndNumberParams{
		StepID:        stepID,
		AttemptNumber: int32(request.AttemptNumber),
	})
	if err == nil {
		if !existing.WorkerID.Valid {
			return nil, errors.New("existing attempt has no worker")
		}
		if err := txQueries.CompleteSchedulingRequest(ctx, pgqueries.CompleteSchedulingRequestParams{
			ID:          claim.ID,
			ClaimID:     claim.ClaimID,
			State:       "DISPATCHED",
			AvailableAt: timestamp(time.Now().UTC()),
		}); err != nil {
			return nil, err
		}
		if err := tx.Commit(ctx); err != nil {
			return nil, err
		}
		existingWorker := uuid.UUID(existing.WorkerID.Bytes)
		_ = dispatcher.waker.WakeWorker(existingWorker.String())
		return &executionports.AttemptRef{
			AttemptID:    existing.ID.String(),
			LeaseVersion: existing.LeaseVersion,
		}, nil
	}
	if !errors.Is(err, pgx.ErrNoRows) {
		return nil, err
	}
	now := time.Now().UTC()
	workerID, generation, err := dispatcher.reserveWorker(ctx, txQueries, stepID, request, now)
	if err != nil {
		if errors.Is(err, errNoEligibleWorker) {
			if commitErr := tx.Commit(ctx); commitErr != nil {
				return nil, commitErr
			}
		}
		return nil, err
	}
	contextRow, err := txQueries.FindStepExecutionAndWorkerNode(ctx, pgqueries.FindStepExecutionAndWorkerNodeParams{
		StepID:   stepID,
		WorkerID: workerID,
	})
	if err != nil {
		return nil, err
	}
	executionID, nodeID := contextRow.ExecutionID, contextRow.NodeID
	if err := appendTimelineTx(ctx, txQueries, timelineRecord{
		ExecutionID:   executionID,
		MilestoneKey:  "step.worker.selected:" + stepID.String() + ":" + fmt.Sprint(request.AttemptNumber),
		Stage:         "SCHEDULING",
		AttemptNumber: request.AttemptNumber,
		State:         "SUCCEEDED",
		Component:     "Scheduler",
		Title:         "已选定 Node 与 Capability 唯一 Worker",

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

知识迁移:选择函数算出结果,事务预留才决定事实。这一分工适用于库存分配、席位分配和租约任务。

5. Worker 执行:本地事实保护进程边界 ​

SSE 下发命令,Worker SQLite 保存 Inbox 与 Attempt 后推进执行。具体 Capability 处理文件,进度与 Checkpoint 提供恢复依据。AccessGrant 按稳定资源引用给出短期读写授权,文件字节由 Worker 直接传输。

能力结果先进入本地事件 Outbox,再由唯一发送器上报。网络失败保留事件重试;本地数据库错误必须独立处理,不能当作暂时断网吞掉。

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

Worker 的依赖与运行状态 backend · c6dbe05c

worker/internal/core/execution/application/runtime.go · 第 48–100 行
符号:Runtime · 核对日期 2026-10-02
来源与提交版本一致

// Runtime 是 Worker 的长期运行应用:注册会话、恢复本地状态、消费命令并可靠上报事件。
type Runtime struct {
	registration       Registration
	control            ControlPlane
	commands           CommandStream
	store              LocalStore
	capabilityExecutor CapabilityExecutor
	ids                executionports.IDGenerator
	observer           executionports.Observation
	now                func() time.Time
	retryInitial       time.Duration
	outboxWake         chan struct{}                 // 合并事件发送唤醒,可靠数据始终保存在本地数据库。
	outboxDrain        chan chan error               // 请求发送器完成一轮发送并返回结果。
	admissionMu        sync.Mutex                    // 串行化本地容量检查与尝试接收。
	activeMu           sync.Mutex                    // 保护正在执行的协程及其取消函数。
	active             map[string]context.CancelFunc // 本进程持有的运行中尝试。
	canceled           map[string]bool               // 记录当前会话已收到的取消。
	draining           bool                          // 当前进程的排空状态,持久标记用于重启恢复。
	resourceInventory  func() executiondomain.NodeInventory
	resourceUsage      func(context.Context) executiondomain.ResourceUsage
}

func New(registration Registration, control ControlPlane, commands CommandStream, store LocalStore, capabilityExecutor CapabilityExecutor, ids executionports.IDGenerator, observer executionports.Observation) (*Runtime, error) {
	if registration.InstanceID == "" || registration.Name == "" || registration.NodeID == "" {
		return nil, executiondomain.Wrap(
			executiondomain.Invalid,
			"Worker 注册信息缺少实例、名称或 Node",
			errors.New("worker instance ID, name, and node ID are required"),
		)
	}
	if registration.Inventory.CPU.CapacityMillis < 1 {
		return nil, executiondomain.Wrap(
			executiondomain.Invalid,
			"Worker CPU 资源采集不可用",
			errors.New("worker CPU inventory is required"),
		)
	}
	if control == nil || commands == nil || store == nil || capabilityExecutor == nil || ids == nil || observer == nil {
		return nil, errors.New("control plane, command stream, local store, capability, ID generator, and observer are required")
	}
	return &Runtime{
		registration:       registration,
		control:            control,
		commands:           commands,
		store:              store,
		capabilityExecutor: capabilityExecutor,
		ids:                ids,
		observer:           observer,
		now:                time.Now,
		retryInitial:       time.Second,
		outboxWake:         make(chan struct{}, 1),
		outboxDrain:        make(chan chan error),
		active:             make(map[string]context.CancelFunc),

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

单发送器及网络重试 backend · c6dbe05c

worker/internal/core/execution/application/outbox.go · 第 34–87 行
符号:sendEvents · 核对日期 2026-10-02
来源与提交版本一致

// sendEvents 独占当前会话的事件发送;能力协程只负责持久化并唤醒发送器。
// 临时网络失败保留事件重试,不取消其他运行中的尝试。
func (runtime *Runtime) sendEvents(ctx context.Context, session Session) error {
	delay := time.Second
	timer := time.NewTimer(0)
	defer timer.Stop()
	for {
		var drained chan error
		select {
		case <-ctx.Done():
			return nil
		case <-runtime.outboxWake:
		case drained = <-runtime.outboxDrain:
		case <-timer.C:
		}
		err := runtime.flushEvents(ctx, session)
		if drained != nil {
			drained <- err
		}
		if err == nil {
			delay = time.Second
			timer.Reset(time.Second)
			continue
		}
		if ctx.Err() != nil {
			return nil
		}

		if !retryableDelivery(err) {
			return err
		}
		runtime.observer.Event(ctx, "event delivery deferred", err)
		// 退避期间不响应生产者唤醒;新增事件仍保留在本地数据库中。
		retry := time.NewTimer(delay)
		select {
		case <-ctx.Done():
			retry.Stop()
			return nil
		case <-retry.C:
		}
		timer.Reset(0)
		delay = min(delay*2, 30*time.Second)
	}
}

// 仅标记远程投递失败;本地数据库错误不能被当作网络故障无限重试。
type deliveryError struct{ error }

func (e *deliveryError) Unwrap() error { return e.error }
func retryableDelivery(err error) bool {
	var delivery *deliveryError
	return errors.As(err, &delivery) && domain.Kind(err) != domain.Unauthenticated && domain.Kind(err) != domain.Forbidden
}

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

资源授权用例 backend · c6dbe05c

server/internal/core/accessgrant/application/service.go · 第 61–158 行
符号:Service.Authorize / issue · 核对日期 2026-10-02
来源与提交版本一致

func (service *Service) Authorize(ctx context.Context, request domain.Authorization) (*domain.Grant, error) {
	if request.AttemptID == "" || request.LeaseVersion < 1 || request.ResourceRef == "" || request.IdempotencyKey == "" ||
		(request.Operation != domain.Download && request.Operation != domain.Upload) ||
		(request.ScopeType == "") != (request.ScopeResourceRef == "") || request.ScopeType != "" && request.Operation != domain.Download {
		return nil, errors.New("invalid resource authorization")
	}
	grantID, err := service.ids.New()
	if err != nil {
		return nil, err
	}
	stored, err := service.repository.ReserveBoundGrant(ctx, request, grantID, service.clock.Now().UTC())
	if err != nil {
		return nil, fmt.Errorf("reserve bound AccessGrant: %w", err)
	}
	return service.issue(ctx, stored, request.AttemptID, request.LeaseVersion, request.Operation, request.IdempotencyKey)
}

func (service *Service) issue(ctx context.Context, stored *domain.StoredGrant, attemptID string, leaseVersion int64, operation domain.Operation, idempotencyKey string) (*domain.Grant, error) {
	secret, err := service.cipher.Open(stored.EncryptedBearer, append([]byte("resource-provider:"), stored.ApplicationAAD...))
	if err != nil {
		return nil, err
	}
	issued, err := service.provider.Issue(ctx, domain.IssueRequest{
		GrantID:        stored.GrantID,
		ApplicationID:  stored.ApplicationID,
		ResourceRef:    stored.ResourceRef,
		AttemptID:      attemptID,
		LeaseVersion:   leaseVersion,
		Operation:      operation,
		IdempotencyKey: idempotencyKey,
		ProviderURL:    stored.ProviderURL,
		BearerToken:    string(secret),
		Scope:          stored.Scope,
	})
	if err != nil {
		return nil, fmt.Errorf("issue AccessGrant: %w", err)
	}
	if err := service.validate(*issued, domain.Renewal{
		GrantID:   stored.GrantID,
		Operation: operation,
	}); err != nil {
		return nil, err
	}
	encoded, err := json.Marshal(issued)
	if err != nil {
		return nil, fmt.Errorf("encode AccessGrant: %w", err)
	}
	ciphertext, err := service.cipher.Seal(encoded, stored.GrantAAD)
	if err != nil {
		return nil, err
	}
	if err := service.repository.Save(ctx, stored.GrantID, ciphertext, issued.ExpiresAt, service.clock.Now().UTC()); err != nil {
		return nil, fmt.Errorf("save AccessGrant: %w", err)
	}
	return issued, nil
}

func (service *Service) Allocate(ctx context.Context, request domain.AllocationRequest) (*domain.Allocation, error) {
	if request.AttemptID == "" || request.LeaseVersion < 1 || request.OutputCollectionRef == "" || request.RelativePath == "" ||
		request.ArtifactType == "" || len(request.ArtifactType) > 100 || request.MediaType == "" || request.MaxBytes < 1 {
		return nil, errors.New("invalid output artifact allocation")
	}
	allocationID, err := service.ids.New()
	if err != nil {
		return nil, err
	}
	stored, err := service.repository.ReserveAllocation(ctx, request, allocationID, service.clock.Now().UTC())
	if err != nil {
		return nil, fmt.Errorf("reserve output artifact allocation: %w", err)
	}
	secret, err := service.cipher.Open(stored.EncryptedSecret, append([]byte("resource-provider:"), stored.ApplicationAAD...))
	if err != nil {
		return nil, err
	}
	allocated, err := service.provider.Allocate(ctx, domain.AllocateProviderRequest{
		StoredAllocation: *stored,
		BearerToken:      string(secret),
	})
	if err != nil {
		return nil, fmt.Errorf("allocate output artifact: %w", err)
	}
	if allocated.ResourceRef == "" || allocated.Grant.Operation != domain.Upload || allocated.Grant.Method != "PUT" || allocated.Grant.MaxBytes < 1 ||
		allocated.Grant.ID == "" || !allocated.Grant.ExpiresAt.After(service.clock.Now()) || !validTemporaryURL(allocated.Grant.URL, service.allowInsecureHTTP) {
		return nil, errors.New("Resource Provider returned an invalid output allocation")
	}
	encoded, err := json.Marshal(&allocated.Grant)
	if err != nil {
		return nil, err
	}
	grantAAD := []byte(allocated.Grant.ID)
	ciphertext, err := service.cipher.Seal(encoded, grantAAD)
	if err != nil {
		return nil, err
	}
	if err := service.repository.CompleteAllocation(ctx, *stored, *allocated, ciphertext, service.clock.Now().UTC()); err != nil {
		return nil, fmt.Errorf("complete output artifact allocation: %w", err)
	}
	return allocated, nil

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

故障窗口:执行产物已写而结果未落库,需要能力的幂等产物身份、Checkpoint 或对账;通用 Runtime 无法自动推断任意外部动作是否完成。

6. 结果与回调:事实转成业务投影 ​

Server 接收 Attempt 事件进入 Inbox,按有效身份、租约与状态条件推进事实,再让工作流继续或结束。最终执行事件产生 Callback delivery,App Demo 验证签名、归属和序列,持久化 Callback Inbox 并更新业务投影。

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

执行状态机及单调进度 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 (

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

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
}

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

Callback 是可靠通知,不是共享数据库写入。业务还需查询补偿,处理缺口及旧执行迟到结果。完成处理单与发布产物仍属于业务规则。

数据与身份对照 ​

边界稳定身份持久证据主要恢复机制
业务创建processing order ID / request key处理单与 command_outbox幂等提交、后台重试
平台创建task ID / execution IDTask、Execution、Event、Outbox请求指纹、事务原子性
消息发布event ID / claim IDoutbox_events租约回收、重投
调度派发Step+AttemptNumberAttempt、worker_commands稳定键、锁与容量复核
Worker 接收command ID / Attempt leaseSQLite Inbox / Attempt去重、恢复扫描
Worker 上报event ID本地 Outbox / Server Inbox重试、原子接收
业务投影event ID / execution ID / sequencecallback_events / 处理单去重、防倒退、查询补偿

这里的表强调持久责任,并非数据库全部字段列表。深入事务与约束时,先读 SQL 源,再读适配器与集成测试。

综合迁移练习 ​

把案例改成视频审核:下载视频、并行抽帧与检测、等待人工确认、发布结果。列出每个拥有者、标识、原子事务和未知结果窗口。

参考设计:业务平台拥有视频与发布规则;执行平台拥有 Run/Step/Attempt;外部检测以稳定请求身份幂等;人工确认有审核身份和耐久等待;产物使用版本化引用;回调与查询更新当前运行投影。发布视频是独立业务动作,不能因某个检测 Step 成功就直接公开。

回到 知识地图,选择仍说不清的机制深入。

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