切换主题
代码快照
以下是课程核对时保存的关键实现,包含路径、符号、提交与片段行号。源码可在手机直接展开阅读。完整文件变化通过内容指纹检查;这些版本不代表线上部署版本。
模型、用例与组合
已核对的实现 · 本地代码快照
实体、值对象与创建不变量 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 包或修改集群。
片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
如何解释证据与维护快照,见 证据使用方法。