Skip to content

契约、校验与错误 ​

问题:能解码 JSON 就代表输入有效吗 ​

请求中的 priority=10000 可以转换成整数,缺失名称也能成为空字符串。这只说明输入可以解码,不说明满足 API Schema 或业务规则。TypeScript 类型更不能约束运行时远端 JSON。

一个清晰边界会依次完成身份验证、结构与范围校验、归属检查和业务决策。内部代码可以依赖已经建立的前置条件,避免每层重复猜输入是否可信。

原理:契约是可检验的协作承诺 ​

资源接口描述对象,显式动作描述领域命令。例如更新名称与取消执行的语义不同;把所有动作塞进通用 PATCH 会模糊允许条件、幂等性和审计。

OpenAPI 描述参数、请求体和响应结构;生成绑定代码减少手写漂移;校验中间件执行 Schema 约束;领域模型验证状态不变量。四者相关但不能相互替代。

最小例子:错误要保留操作语义 ​

json
{
  "type": "about:blank",
  "title": "Conflict",
  "status": 409,
  "detail": "申请已经审批,无法再次修改",
  "request_id": "request-example"
}

这是教学 Problem Details 示例。客户端按 status 和稳定分类决策,detail 给人阅读,不用解析错误句子来判断是否重试。内部 SQL 错误、凭据和远端响应体不应直接暴露。

方案比较 ​

手写 DTO 灵活,需人工维护服务端与客户端一致性;生成类型降低重复,但源契约仍需评审和运行时校验。通用 Resource 袋方便管理列表,却不适合作为核心业务协作对象,否则关键字段隐藏在 JSON map 中。

兼容性包括字段、枚举、行为和错误语义。新增枚举可能破坏穷举客户端;可选字段变必填可能破坏旧调用方;仅保持 URL 不变不代表兼容。

真实案例:三种校验分工 ​

创建 Task 应用校验输入、回调、身份与幂等键,并根据发布的 Workflow Schema 验证输入。领域 NewTask 进一步保护名称、优先级和创建不变量。HTTP 组合层还接入 OpenAPI validator;生成 strict handler 的解码并不独自覆盖完整 Schema 校验。

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

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

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

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

OpenAPI 与 SQL 是生成输入,contracts/gen 和 sqlc 输出通过生成目标维护。学习时先读源与手写适配器,排查具体映射时再看生成代码。

失败与边界:错误类型影响重试 ​

结构无效应尽快失败;冲突可能需刷新状态;依赖不可用可以重试,但必须考虑幂等。401 需要重新建立身份,403 不应自动无限重试。超时可能对应未知结果,需稳定请求键和查询补偿。

校验之外还需请求体大小、并发与速率限制。即使格式合法,也不能允许任意规模的输入耗尽系统。服务端分页与客户端遍历必须约定边界。

迁移练习与参考答案 ​

练习:导出 API 要新增“取消”。用 PATCH status=CANCELED 还是 POST /exports/{id}/cancel?

参考答案:取消是有前置条件、外部执行和中间状态的领域动作,显式动作更清楚。返回取消请求已接受的状态,最终结果另查;加入幂等键、版本或状态检查、身份归属和审计。不能让客户端直接选择所有终态。

继续阅读:身份与安全边界、测试与故障验证。

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