refactor: split container executor validation states

This commit is contained in:
2026-08-22 15:59:20 +08:00
parent 68d067d744
commit 67214dd858
+68 -34
View File
@@ -158,45 +158,54 @@ func (e *Executor) run(ctx context.Context, transactionID string, request Reques
if err != nil { if err != nil {
return err return err
} }
switch record.State { done, err := e.advanceState(ctx, transactionID, record.State, request)
case transaction.StateCreated: if err != nil {
if _, err := e.store.Transition(ctx, transactionID, transaction.StateValidating, "backend container validation started"); err != nil { return err
return err }
} if done {
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:
return nil 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 阶段的错误是否可恢复。 // failUnlessRecoverable 判断 prepare 阶段的错误是否可恢复。
// 若 cause 是 UncertainStepError 或 ErrStepConflict,则直接原样返回(保留不确定性以便重放恢复); // 若 cause 是 UncertainStepError 或 ErrStepConflict,则直接原样返回(保留不确定性以便重放恢复);
// 否则将事务标记为 StateFailed 并合并返回 cause 与状态迁移错误。 // 否则将事务标记为 StateFailed 并合并返回 cause 与状态迁移错误。
@@ -321,6 +330,19 @@ func (e *Executor) streamContainerLogs(ctx context.Context, name string, report
// validateRequest 对请求字段做静态校验,确保所有取值精确且自洽。 // validateRequest 对请求字段做静态校验,确保所有取值精确且自洽。
// 任一字段不符合要求时返回描述性错误。 // 任一字段不符合要求时返回描述性错误。
func validateRequest(request Request) error { 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 { switch request.ImageAcquisition {
case ImageAcquisitionLoad: case ImageAcquisitionLoad:
if !filepath.IsAbs(request.ArchivePath) { if !filepath.IsAbs(request.ArchivePath) {
@@ -337,6 +359,10 @@ func validateRequest(request Request) error {
default: default:
return fmt.Errorf("unsupported image acquisition: %q", request.ImageAcquisition) 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 { if request.ImageReference == "" || strings.TrimSpace(request.ImageReference) != request.ImageReference {
return errors.New("exact image reference is required") 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 { if request.ContainerName == "" || strings.TrimSpace(request.ContainerName) != request.ContainerName {
return errors.New("exact container name is required") 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 { 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) 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 { if err := validateRestartPolicy(request.RestartPolicy); err != nil {
return err return err
} }
return nil
}
func validateHealthEndpoint(request Request) error {
path := request.HealthPath path := request.HealthPath
if path == "" { if path == "" {
path = healthPath path = healthPath