From 67214dd8588b8467d8e82e96ecd7256417108da5 Mon Sep 17 00:00:00 2001 From: Zhan Ziyang Date: Sat, 22 Aug 2026 15:59:20 +0800 Subject: [PATCH] refactor: split container executor validation states --- internal/backendexecutor/executor.go | 102 ++++++++++++++++++--------- 1 file changed, 68 insertions(+), 34 deletions(-) diff --git a/internal/backendexecutor/executor.go b/internal/backendexecutor/executor.go index c7dd2e0..51fbc00 100644 --- a/internal/backendexecutor/executor.go +++ b/internal/backendexecutor/executor.go @@ -158,45 +158,54 @@ func (e *Executor) run(ctx context.Context, transactionID string, request Reques if err != nil { return err } - switch record.State { - case transaction.StateCreated: - if _, err := e.store.Transition(ctx, transactionID, transaction.StateValidating, "backend container validation started"); err != nil { - return err - } - case transaction.StateValidating: - if err := e.validate(ctx, request); err != nil { - _, transitionErr := e.store.Transition(ctx, transactionID, transaction.StateFailed, err.Error()) - return errors.Join(err, transitionErr) - } - if _, err := e.store.Transition(ctx, transactionID, transaction.StatePrepared, "backend container inputs validated"); err != nil { - return err - } - case transaction.StatePrepared: - if err := e.prepare(ctx, transactionID, request); err != nil { - return e.failUnlessRecoverable(ctx, transactionID, err) - } - if _, err := e.store.Transition(ctx, transactionID, transaction.StateStarting, "backend container prepared"); err != nil { - return err - } - case transaction.StateStarting: - // 重放 PREPARED 步骤会核对持久化意图,阻止恢复时换入另一组请求参数。 - if err := e.prepare(ctx, transactionID, request); err != nil { - return err - } - if err := e.startAndCheck(ctx, transactionID, request); err != nil { - return err - } - if _, err := e.store.Transition(ctx, transactionID, transaction.StateSwitching, "backend container is healthy"); err != nil { - return err - } - case transaction.StateSwitching: + done, err := e.advanceState(ctx, transactionID, record.State, request) + if err != nil { + return err + } + if done { return nil - default: - return fmt.Errorf("backend container executor cannot run transaction %s in state %s", transactionID, record.State) } } } +func (e *Executor) advanceState(ctx context.Context, transactionID string, state transaction.State, request Request) (bool, error) { + switch state { + case transaction.StateCreated: + _, err := e.store.Transition(ctx, transactionID, transaction.StateValidating, "backend container validation started") + return false, err + case transaction.StateValidating: + return false, e.validateAndPrepare(ctx, transactionID, request) + case transaction.StatePrepared: + if err := e.prepare(ctx, transactionID, request); err != nil { + return false, e.failUnlessRecoverable(ctx, transactionID, err) + } + _, err := e.store.Transition(ctx, transactionID, transaction.StateStarting, "backend container prepared") + return false, err + case transaction.StateStarting: + if err := e.prepare(ctx, transactionID, request); err != nil { + return false, err + } + if err := e.startAndCheck(ctx, transactionID, request); err != nil { + return false, err + } + _, err := e.store.Transition(ctx, transactionID, transaction.StateSwitching, "backend container is healthy") + return false, err + case transaction.StateSwitching: + return true, nil + default: + return false, fmt.Errorf("backend container executor cannot run transaction %s in state %s", transactionID, state) + } +} + +func (e *Executor) validateAndPrepare(ctx context.Context, transactionID string, request Request) error { + if err := e.validate(ctx, request); err != nil { + _, transitionErr := e.store.Transition(ctx, transactionID, transaction.StateFailed, err.Error()) + return errors.Join(err, transitionErr) + } + _, err := e.store.Transition(ctx, transactionID, transaction.StatePrepared, "backend container inputs validated") + return err +} + // failUnlessRecoverable 判断 prepare 阶段的错误是否可恢复。 // 若 cause 是 UncertainStepError 或 ErrStepConflict,则直接原样返回(保留不确定性以便重放恢复); // 否则将事务标记为 StateFailed 并合并返回 cause 与状态迁移错误。 @@ -321,6 +330,19 @@ func (e *Executor) streamContainerLogs(ctx context.Context, name string, report // validateRequest 对请求字段做静态校验,确保所有取值精确且自洽。 // 任一字段不符合要求时返回描述性错误。 func validateRequest(request Request) error { + if err := validateImageAcquisition(request); err != nil { + return err + } + if err := validateImageIdentity(request); err != nil { + return err + } + if err := validateContainerRuntime(request); err != nil { + return err + } + return validateHealthEndpoint(request) +} + +func validateImageAcquisition(request Request) error { switch request.ImageAcquisition { case ImageAcquisitionLoad: if !filepath.IsAbs(request.ArchivePath) { @@ -337,6 +359,10 @@ func validateRequest(request Request) error { default: return fmt.Errorf("unsupported image acquisition: %q", request.ImageAcquisition) } + return nil +} + +func validateImageIdentity(request Request) error { if request.ImageReference == "" || strings.TrimSpace(request.ImageReference) != request.ImageReference { return errors.New("exact image reference is required") } @@ -349,6 +375,10 @@ func validateRequest(request Request) error { if request.ContainerName == "" || strings.TrimSpace(request.ContainerName) != request.ContainerName { return errors.New("exact container name is required") } + return nil +} + +func validateContainerRuntime(request Request) error { if request.Port != 8080 && request.Port != 8081 && request.Port != 18910 && request.Port != 28910 { return fmt.Errorf("container port must be 8080, 8081, 18910, or 28910: %d", request.Port) } @@ -372,6 +402,10 @@ func validateRequest(request Request) error { if err := validateRestartPolicy(request.RestartPolicy); err != nil { return err } + return nil +} + +func validateHealthEndpoint(request Request) error { path := request.HealthPath if path == "" { path = healthPath