Skip to content

代码快照 ​

以下是课程核对时保存的关键实现,包含路径、符号、提交与片段行号。源码可在手机直接展开阅读。完整文件变化通过内容指纹检查;这些版本不代表线上部署版本。

模型、用例与组合 ​

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

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

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

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

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

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

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

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

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

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

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

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

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

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

server/internal/core/execution/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)
}

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

组合根及资源创建 backend · c6dbe05c

server/internal/bootstrap/runtime.go · 第 68–167 行
符号:run · 核对日期 2026-10-02
来源与提交版本一致

func run(ctx context.Context, serverConfig config.Server, logger *slog.Logger) (returnedErr error) {
	if logger == nil {
		return errors.New("server logger is required")
	}
	ctx, cancel := context.WithCancel(ctx)
	defer cancel()
	var startupCtx context.Context
	var startupCancel context.CancelFunc
	ctx = operation.Start(ctx, operation.Operation{
		Name: "server.run",
	})
	defer func() { returnedErr = operation.WithFailure(ctx, returnedErr) }()
	startupCtx, startupCancel = context.WithTimeout(ctx, 30*time.Second)
	defer startupCancel()
	poolConfig, err := pgxpool.ParseConfig(serverConfig.DatabaseURL)
	if err != nil {
		return fault.Wrap(fault.Invalid, "postgresql.configure", err)
	}
	poolConfig.ConnConfig.Tracer = observability.NewPostgreSQLQueryLogger(logger.With("component", "postgresql"), serverConfig.SlowQueryThreshold)
	pool, err := pgxpool.NewWithConfig(startupCtx, poolConfig)
	if err != nil {
		return fault.Wrap(fault.Unavailable, "postgresql.connect", err)
	}
	defer pool.Close()
	if err := pool.Ping(startupCtx); err != nil {
		return fault.Wrap(fault.Unavailable, "postgresql.ping", err)
	}
	var baseline string
	if err := pool.QueryRow(startupCtx, "SELECT version FROM tdp_schema_baseline").Scan(&baseline); err != nil || baseline != "application-v1" {
		return errors.New("TDP application-v1 database required; initialize a fresh project database")
	}
	bus, err := natsmessaging.Connect(serverConfig.NATSURL, "tdp-server", 10*time.Second, logger.With("component", "nats"))
	if err != nil {
		return fault.Wrap(fault.Unavailable, "nats.connect", err)
	}
	defer func() {
		if err := bus.Close(); err != nil {
			logger.ErrorContext(ctx, "close NATS", "error", err)
		}
	}()
	if err := bus.EnsureStreams(startupCtx); err != nil {
		return fault.Wrap(fault.Unavailable, "nats.ensure_streams", err)
	}
	cryptoBox, err := cryptobox.NewBase64(serverConfig.DataEncryptionKey)
	if err != nil {
		return err
	}
	definitionValidator, err := definitionvalidation.New()
	if err != nil {
		return err
	}
	grantOptions := []accessgrantapp.Option{}
	if serverConfig.AllowInsecureResourceURLs {
		logger.WarnContext(ctx, "insecure HTTP resource URLs are enabled; use only in local E2E")
		grantOptions = append(grantOptions, accessgrantapp.WithInsecureHTTPForLocalDevelopment())
	}
	grantRepository, err := accessgrantpostgres.NewRepository(pool)
	if err != nil {
		return err
	}
	resourceProvider, err := resourceprovider.New(15 * time.Second)
	if err != nil {
		return err
	}
	grantModule, err := accessgrantapp.NewModule(accessgrantapp.ModuleDependencies{
		Repository: grantRepository,
		Provider:   resourceProvider,
		Cipher:     cryptoBox,
		Clock:      clock.System{},
		IDs:        idgen.StringUUIDv7{},
		Options:    grantOptions,
	})
	if err != nil {
		return err
	}
	grantIssuer, err := workercontroladapter.NewAccessGrantIssuer(grantModule.Service)
	if err != nil {
		return err
	}
	workerStore, err := postgresadapter.NewWorkerControlStore(pool)
	if err != nil {
		return err
	}
	workerControlModule, err := workercontrol.NewModule(workercontrol.ModuleDependencies{
		Store:               workerStore,
		IDs:                 idgen.StringUUIDv7{},
		HeartbeatInterval:   serverConfig.WorkerControl.HeartbeatInterval,
		LeaseDuration:       serverConfig.WorkerControl.LeaseDuration,
		CommandPollInterval: serverConfig.WorkerControl.CommandPollInterval,
		WakeSubscriber:      bus,
		GrantIssuer:         grantIssuer,
		CapabilityValidator: definitionValidator,
	})
	if err != nil {
		return err
	}
	taskRepositoryOptions := []postgresadapter.TaskRepositoryOption{}
	if serverConfig.AllowInsecureCallbackURLs {
		taskRepositoryOptions = append(taskRepositoryOptions, postgresadapter.WithInsecureCallbacksForLocalDevelopment())
	}

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

Capability 可替换执行契约 backend · c6dbe05c

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

// CapabilityExecutor 将通用执行命令委托给具体 Capability。
type CapabilityExecutor interface {
	// Execute 根据命令中的能力名称和版本执行,并通过 reporter 上报进度。
	Execute(context.Context, domain.Command, ProgressReporter) (*domain.AttemptEvent, error)
}

// Progress 是能力执行过程中上报给控制面的结构化进度。
type Progress struct {
	Stage     string `json:"stage"`        // 当前处理阶段名称。
	Status    string `json:"stage_status"` // 阶段状态。
	Progress  int    `json:"progress"`     // 整体 0 到 100 的进度。
	Completed int    `json:"completed"`    // 当前阶段已完成条目数。
	Total     int    `json:"total"`        // 当前阶段总条目数。
	Failed    int    `json:"failed"`       // 当前阶段失败条目数。
}

// ProgressReporter 将能力进度转换为本地检查点和控制面尝试心跳。
type ProgressReporter func(context.Context, Progress) error

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

可靠性、调度与执行 ​

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

发布、认领与退避 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);

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

硬约束筛选与确定性评分 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",

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

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

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

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

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

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

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

前端与业务边界 ​

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

CSRF、错误与幂等请求 frontend · b0308f40

src/api/client.ts · 第 1–60 行
符号:http / write · 核对日期 2026-10-02
来源与提交版本一致

import axios from "axios";
import { user } from "@/auth";
import type { Resource, ResourcePage } from "./types";
export const http = axios.create({ baseURL: "/api/v1", timeout: 20000 });
http.interceptors.request.use((c) => {
  if (user.value && !["get", "head", "options"].includes(c.method || "get"))
    c.headers["X-CSRF-Token"] = user.value.csrf_token;
  return c;
});
export function errorText(e: unknown): string {
  if (axios.isAxiosError(e)) {
    const p = e.response?.data;
    return (
      (p?.detail || p?.title || e.message) +
      (p?.request_id ? `(请求 ${p.request_id})` : "")
    );
  }
  return e instanceof Error ? e.message : String(e);
}
http.interceptors.response.use(
  (r) => r,
  async (e) => {
    if (e.response?.status === 401) {
      user.value = null;
      if (!location.pathname.endsWith("/login")) location.assign("/login");
    }
    return Promise.reject(e);
  },
);
export async function get<T>(
  path: string,
  params?: object,
  signal?: AbortSignal,
) {
  return (await http.get<T>(path, { params, signal })).data;
}
export async function write<T = Resource>(
  path: string,
  body: unknown,
  key: string,
  method: "post" | "put" | "patch" | "delete" = "post",
) {
  return (
    await http.request<T>({
      url: path,
      method,
      data: body,
      headers: { "Idempotency-Key": key },
    })
  ).data;
}
export async function choices(path: string) {
  const items: Resource[] = [];
  for (let page = 1; ; page++) {
    const p = await get<ResourcePage>(path, { page, page_size: 200 });
    items.push(...p.items);
    if (items.length >= p.total || p.items.length === 0) return items;
  }
}

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

会话与登录状态 frontend · b0308f40

src/auth.ts · 第 2–46 行
符号:authReady · 核对日期 2026-10-02
来源与提交版本一致

export interface AdminSession {
  principal_id: string;
  email: string;
  csrf_token: string;
  expires_at: string;
}
export const user = ref<AdminSession | null>(null);
export const authError = ref("");
export const isAuthenticated = computed(
  () => !!user.value && Date.parse(user.value.expires_at) > Date.now(),
);
export const authReady = fetch("/api/v1/auth/session", {
  credentials: "same-origin",
  cache: "no-store",
})
  .then(async (r) => {
    if (r.status === 401) {
      user.value = null;
      return;
    }
    if (!r.ok) throw new Error("无法读取登录会话,请检查服务状态");
    user.value = await r.json();
  })
  .catch((e) => {
    authError.value = e instanceof Error ? e.message : String(e);
    user.value = null;
  });
export function login(returnTo = "/overview") {
  const target = new URL("/api/v1/auth/login", location.origin);
  target.searchParams.set("return_to", returnTo);
  location.assign(target.pathname + target.search);
}
export async function logout() {
  if (user.value) {
    const r = await fetch("/api/v1/auth/logout", {
      method: "POST",
      credentials: "same-origin",
      headers: { "X-CSRF-Token": user.value.csrf_token },
    });
    if (!r.ok) throw new Error("退出失败,请重试");
  }
  user.value = null;
  location.assign("/login");
}

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

按功能组合路由 frontend · b0308f40

src/router/index.ts · 第 1–50 行
符号:router · 核对日期 2026-10-02
来源与提交版本一致

import { createRouter, createWebHistory } from "vue-router";
import { applicationRoutes } from "@/features/applications/routes";
import { auditRoutes } from "@/features/audit/routes";
import { authRoutes } from "@/features/auth/routes";
import { callbackRoutes } from "@/features/callbacks/routes";
import { definitionRoutes } from "@/features/definitions/routes";
import { fleetRoutes } from "@/features/fleet/routes";
import { overviewRoutes } from "@/features/overview/routes";
import { systemRoutes } from "@/features/system/routes";
import { taskRoutes } from "@/features/tasks/routes";
import { authReady, user } from "@/auth";

const router = createRouter({
  history: createWebHistory("/"),
  routes: [
    ...authRoutes,
    {
      path: "/",
      component: () => import("@/layouts/ConsoleLayout.vue"),
      children: [
        { path: "", redirect: "/overview" },
        ...overviewRoutes,
        ...systemRoutes,
        ...taskRoutes,
        ...fleetRoutes,
        ...applicationRoutes,
        ...definitionRoutes,
        ...callbackRoutes,
        ...auditRoutes,
      ],
    },
    {
      path: "/:pathMatch(.*)*",
      component: () => import("@/shared/pages/NotFound.vue"),
    },
  ],
});

router.beforeEach(async (to) => {
  if (to.meta.public) return;
  await authReady;
  if (
    to.path !== "/login" &&
    (!user.value || Date.parse(user.value.expires_at) <= Date.now())
  )
    return { path: "/login", query: { returnTo: to.fullPath } };
});

export default router;

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

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

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

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
}

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

工程资料 ​

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

生成与工程门禁 backend · c6dbe05c

Makefile · 第 37–83 行
符号:generate / check · 核对日期 2026-10-02
来源与提交版本一致

generate:
	$(GO) tool oapi-codegen -config contracts/oapi-codegen/public.yaml contracts/openapi/public/v1/openapi.yaml
	$(GO) tool oapi-codegen -config contracts/oapi-codegen/worker.yaml contracts/openapi/worker/v1/openapi.yaml
	$(GO) tool oapi-codegen -config contracts/oapi-codegen/management.yaml contracts/openapi/management/v1/openapi.yaml
	$(GO) tool oapi-codegen -config contracts/oapi-codegen/callback.yaml contracts/openapi/callback/v1/openapi.yaml
	$(GO) tool oapi-codegen -config contracts/oapi-codegen/resource-provider.yaml contracts/openapi/resource-provider/v1/openapi.yaml
	$(GO) tool sqlc generate
	$(GO) run ./tools/sqlcpostprocess

migration-validate:
	$(GO) tool goose -dir server/db/migrations validate
	$(GO) tool goose -dir worker/db/migrations validate

generate-check: generate
	git diff --exit-code -- contracts/gen/go server/internal/adapters/outbound/postgres/pgqueries worker/internal/adapters/outbound/sqlite/sqlitequeries

schema-check:
	$(GO) test ./contracts/...

config-check:
	TDP_SERVER_OIDC_CLIENT_SECRET=validation TDP_SERVER_DATABASE_URL=postgres://tdp:validation@localhost/tdp TDP_SERVER_NATS_URL=nats://localhost:4222 TDP_SERVER_DATA_ENCRYPTION_KEY=BwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwcHBwc= $(GO) run ./server/cmd/server config validate --config configs/examples/server.yaml
	$(GO) run ./worker/cmd/worker config validate --config configs/examples/worker.yaml

architecture-check:
	$(GO) run ./tools/archcheck

api-compat:
	@test -n "$(BASE_OPENAPI_DIR)" || (echo "BASE_OPENAPI_DIR is required"; exit 2)
	$(GO) tool oasdiff breaking $(BASE_OPENAPI_DIR)/public/v1/openapi.yaml contracts/openapi/public/v1/openapi.yaml
	$(GO) tool oasdiff breaking $(BASE_OPENAPI_DIR)/worker/v1/openapi.yaml contracts/openapi/worker/v1/openapi.yaml

vuln:
	$(GO) tool govulncheck ./...

license-check:
	$(GO) tool go-licenses check ./server/cmd/server ./worker/cmd/worker \
		--ignore=tdp \
		--ignore=github.com/nexus-rpc/nexus-proto-annotations/go/nexusannotations/v1

check: fmt-check vet test schema-check config-check architecture-check migration-validate

clean:
	$(GO) clean ./...
	rm -f $(BIN_DIR)/server $(BIN_DIR)/worker
	rm -f $(BIN_DIR)/server-linux-amd64 $(BIN_DIR)/worker-linux-amd64
	rm -f $(BIN_DIR)/worker-linux-arm64 $(BIN_DIR)/worker-windows-amd64.exe

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

本地栈与验证入口 workbench · 514f5a00

Makefile · 第 19–55 行
符号:check / up / e2e · 核对日期 2026-10-02
来源与提交版本一致

check:
	$(COMPOSE) -f compose.yaml config --quiet
	$(COMPOSE) -f compose.yaml -f compose.dev.yaml config --quiet
	$(COMPOSE) -f compose.yaml -f compose.dev.yaml -f compose.demo.yaml config --quiet
	$(COMPOSE) -f compose.yaml -f compose.e2e.yaml config --quiet
	$(COMPOSE) -f compose.yaml -f compose.demo.yaml config --quiet
	python3 scripts/check-doc-links.py

docs-check:
	python3 scripts/check-doc-links.py

ps:
	$(LOCAL_COMPOSE) ps -a

up:
	$(LOCAL_COMPOSE) up -d --build

down:
	$(LOCAL_COMPOSE) down

logs:
	$(LOCAL_COMPOSE) logs -f

e2e:
	TDP_E2E_PROJECT=$(E2E_PROJECT) ./scripts/e2e.sh

e2e-down:
	$(COMPOSE) -p $(E2E_PROJECT) -f compose.yaml -f compose.e2e.yaml $(E2E_EXTRA_COMPOSE) --profile test down

fault-test:
	TDP_E2E_PROJECT=$(E2E_PROJECT) ./scripts/fault-test.sh

demo-up:
	$(LOCAL_COMPOSE) -f compose.demo.yaml up -d --build

demo-down:
	$(LOCAL_COMPOSE) -f compose.demo.yaml down

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

日志关联与当前观测边界 backend · c6dbe05c

docs/implementation/logging.md · 第 27–68 行
符号:日志记录责任与关联 · 核对日期 2026-10-02
来源与提交版本一致

## 记录责任

- Core/Repository/Client 分类和包装错误;errors.Is/As 可继续检查原始 cause。
- HTTP 访问日志为 INFO;未预期错误由响应边界额外记录一次 ERROR,响应只返回安全 Problem Details。
- 认证无效为 401,禁止访问为 403,依赖不可用为 503,未知内部故障为 500。
- 消息处理失败记录投递次数、durable 和事件 ID,保持 NAK 延迟及原有最大投递次数;最终一次失败为 ERROR。
- 后台重试为 WARN,重试耗尽或进程退出为 ERROR;空轮询、正常关闭和预期取消不写错误日志。
- Worker 保留会话就绪、Attempt 开始/终态及失败决策日志;成功命令 ACK 降为 DEBUG,不为正常心跳新增 INFO 事件。HTTP 访问日志仍遵循上述统一规则。
- Worker Capability 已转换为 FAILED 事件的错误在转换处记录;事件上报失败由会话重试边界记录。
- timeline/审计辅助写入失败为 WARN,业务持久化错误继续返回。
- Temporal 使用 SDK 的 slog 适配器及 replay-aware Workflow logger;Activity 创建本地操作上下文。

## 关联与脱敏

基础字段为 service/version/revision/instance_id。操作日志附带 operation、trace_id,
以及当时已知的 request_id、worker_id、node_id、command_id、attempt_id、execution_id、event_id。
耗时统一为 duration_ms。异步错误通过 Failure 只保留来源操作的关联元数据,不持有整个 context。
独立请求、消息投递、轮询轮次、Worker 会话及命令创建本地 ID;命令启动的 Attempt、事件上报、
心跳及 SSE 重连复用所属范围的 ID。重试重新进入独立执行边界时才创建新 ID,跨进程使用业务 ID。
HTTP 的 request_id 与 trace_id 字段继续保留,以免改变现有日志及 Problem Details 消费方式。
这里的 trace_id 是日志关联 ID,不代表可以查询一棵 span 树。

PostgreSQL 和 SQLite 查询日志通常为 DEBUG,耗时达到阈值时为 WARN;
查询日志与边界错误日志分工明确,不在 SQL 层再输出 ERROR。
PostgreSQL 的 `sql` 字段记录脱敏后的语句,Worker SQLite 的 `sql` 字段记录稳定的 sqlc 查询名;
SQLite 包含事务查询以及 Scan/Rows 错误;启动迁移由启动边界诊断。
参数、外部响应体、签名 URL、密码、令牌及 PostgreSQL Detail 不进入日志。
未知错误只输出分类和类型,数据库错误保留 SQLSTATE;需要更详细诊断时,
在适配器补充安全的 Fault.Operation,不能直接开放任意 err.Error()。

在 Workbench 的 Loki 中可用以下方式关联同一进程内的操作:

```logql
{compose_project="tdp-workbench"} | json | trace_id="<trace-id>"
```

跨 Server/Worker 查询使用 attempt_id、command_id 或 event_id。本改动不需要新建
数据库字段或部署 Trace 收集服务。

## 使用与日志示例

临时开启 SQL(下一次启动生效):

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

跨项目端到端资料 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。

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

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

设计要求中的系统边界 design · fb16fc7d

README.md · 第 14–33 行
符号:核心边界 · 核对日期 2026-10-02
来源与提交版本一致

核心边界:

- 业务系统自身的任务与 TDP Task 是不同对象,通过稳定的不透明引用和 HTTP 契约关联;TDP 内部采用 `Task -> Execution -> Step -> Attempt`;
- TDP 必须将异步执行的进展和最终状态回调给业务平台,业务平台同时可通过查询 API 进行补偿;
- TDP 不理解资产、课程、妆容、镜头等具体行业模型;
- TDP Server 与 Worker 采用同一代码仓库中的两个独立应用,分别构建、发布和部署,只通过版本化协议通信;
- TDP 自有 HTTP 业务接口统一使用 `/api/v1/...`,Public、Worker、Management 边界由独立 Listener、Service 或域名区分;
- 自有接口采用“RESTful 资源 + 显式领域动作”:资源管理遵循 RESTful 风格,注册、心跳、确认、取消等命令优先保证语义清晰,并具备幂等与审计约束;
- Hroxy 与 Omni Server 处于相同架构层级,都是 Worker 的调用或调度对象;它们不替代 Worker、不注册到 Node,也不直接参与 TDP 协议;
- TDP 只调度 Worker,具体 Capability Worker 负责选择和调用可达的下游执行服务、保存其外部引用并把结果转换为 TDP Attempt 事件;并非每个 Worker 都需要下游服务;
- 原始资源和执行产物归业务平台所有,TDP 不持久化业务文件;
- 所有项目统一采用 Go、六边形架构和 DDD;
- 实现顺序优先打通核心纵向切片,但第一阶段交付完整的集群、可靠性、安全、监控和治理能力,不以最小实现名义删减;
- 第三方包优先采用当前仍受维护的最新稳定版本,并固定实际构建版本;不无依据沿用陈旧、停更或存在已知安全风险的依赖;
- 业务平台和 TDP 的服务端关系数据使用 PostgreSQL,Worker 本地结构化状态使用 SQLite;分布式协调、队列、OpenAPI/Swagger、结构化日志、本地关联、Server 指标和健康检查同属当前程序实现基线,集中指标、告警与分布式 Trace 仍属后续建设。

## 总体架构

```mermaid
flowchart LR

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

发布仓库与运维职责 deploy · cfb60638

README.md · 第 1–15 行
符号:部署仓库说明 · 核对日期 2026-10-02
来源与提交版本一致

# TDP 发布

本仓库 [deploy-k3s](https://cnb.cool/iosoon/tdp/deploy-k3s) 是 Workbench 的可选 `deploy/` 子模块,独立维护两条并列发布流程:

| 目录 | 范围 | 手册 |
| --- | --- | --- |
| `k3s/` | Server、管理前端、App Demo 的镜像、K3s 清单与部署 | [K3s 发布](k3s/README.md) |
| `worker/` | Worker 五平台包、MinIO 发行与目标机升级 | [Worker 发布](worker/README.md) |

版本审计及最近生产基线见 [DEPLOYMENT.md](DEPLOYMENT.md)。共享基础设施由 `io-k3s` 管理;本仓库不管理 Workbench 业务源码,也不为日常代码开发增加部署步骤。

在 Workbench 根目录运行 `make deploy-init`,部署维护者随后在 `deploy/` 中切换到 `main` 并单独提交、推送。普通开发只运行 Workbench 的 `make init`,不会初始化本仓库。独立检出时将 `WORKBENCH_ROOT` 指向有 `backend/`、`frontend/`、`app-demo/` 的源码目录。

本地检查:`make check`。任何实际发布都需要用户明确确认版本、组件和目标;改动发布工具不授权构建推送镜像、公开 Worker 包或修改集群。

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

如何解释证据与维护快照,见 证据使用方法。

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