切换主题
一条任务的旅程
这个案例用于检验课程中的知识。业务是上传 ZIP,清理不需要的文件,再给图片加水印。关注每段链路的所有权、原子边界和恢复证据,而不是记忆项目路径。
总体链路
100%
业务处理单
Command Outbox
→Command Outbox
TDP Task / Execution
Event / Outbox
Event / Outbox
NATS → Temporal DAG
→Step / Attempt
Worker Command
Worker Command
Worker SQLite Inbox
Capability / Event Outbox
→Capability / Event Outbox
Server Attempt Inbox
状态推进 / Signal
状态推进 / 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 ID | Task、Execution、Event、Outbox | 请求指纹、事务原子性 |
| 消息发布 | event ID / claim ID | outbox_events | 租约回收、重投 |
| 调度派发 | Step+AttemptNumber | Attempt、worker_commands | 稳定键、锁与容量复核 |
| Worker 接收 | command ID / Attempt lease | SQLite Inbox / Attempt | 去重、恢复扫描 |
| Worker 上报 | event ID | 本地 Outbox / Server Inbox | 重试、原子接收 |
| 业务投影 | event ID / execution ID / sequence | callback_events / 处理单 | 去重、防倒退、查询补偿 |
这里的表强调持久责任,并非数据库全部字段列表。深入事务与约束时,先读 SQL 源,再读适配器与集成测试。
综合迁移练习
把案例改成视频审核:下载视频、并行抽帧与检测、等待人工确认、发布结果。列出每个拥有者、标识、原子事务和未知结果窗口。
参考设计:业务平台拥有视频与发布规则;执行平台拥有 Run/Step/Attempt;外部检测以稳定请求身份幂等;人工确认有审核身份和耐久等待;产物使用版本化引用;回调与查询更新当前运行投影。发布视频是独立业务动作,不能因某个检测 Step 成功就直接公开。
回到 知识地图,选择仍说不清的机制深入。