diff --git a/internal/backendupdate/container.go b/internal/backendupdate/container.go index 45a9fe1..abbb59c 100644 --- a/internal/backendupdate/container.go +++ b/internal/backendupdate/container.go @@ -517,53 +517,15 @@ func (u *Updater) switchAndCommitContainer(ctx context.Context, transactionID st } switch record.State { case transaction.StateSwitching: - reportProgress(report, Progress{TransactionID: transactionID, State: record.State, Message: fmt.Sprintf("Switching host Nginx backend traffic to port %d", request.TargetPort)}) - operation := &gatewayOperation{controller: u.gateway, before: before, after: after, receiptPath: request.GatewayReceiptPath} - if _, err := u.coordinator.ExecuteStep(ctx, transactionID, containerGatewaySwitchIntent(request), operation); err != nil { - return u.rollbackContainer(ctx, transactionID, request, before, after, err) - } - if _, err := u.store.Transition(ctx, transactionID, transaction.StateVerifying, "host Nginx now routes backend traffic to the healthy container slot"); err != nil { + if err := u.switchContainerTraffic(ctx, transactionID, request, before, after, report); err != nil { return err } case transaction.StateVerifying: - message := "container backend switch verified" - if request.PreviousContainer != "" { - message += "; previous container draining" - } - if _, err := u.store.Transition(ctx, transactionID, transaction.StateDraining, message); err != nil { + if err := u.markContainerDraining(ctx, transactionID, request); err != nil { return err } case transaction.StateDraining: - if request.PreviousContainer != "" { - reportProgress(report, Progress{TransactionID: transactionID, State: record.State, Message: fmt.Sprintf("Draining previous backend container for %s", u.drain)}) - if err := waitContext(ctx, u.drain); err != nil { - return err - } - operation := &containerStopOperation{engine: u.engine, name: request.PreviousContainer, stopGraceSeconds: u.containerStopGraceSeconds} - if _, err := u.coordinator.ExecuteStep(ctx, transactionID, stopPreviousContainerIntent(request), operation); err != nil { - return err - } - } - target, err := u.engine.InspectContainer(ctx, request.TargetContainer) - if err != nil { - return fmt.Errorf("inspect target backend container before commit: %w", err) - } - if !target.Running || target.Dead || target.ID == "" { - return fmt.Errorf("target backend container %s is not running with an exact identity", request.TargetContainer) - } - if target.ImageReference != request.ImmutableReference && target.ImageReference != request.ImageReference { - return fmt.Errorf("target backend container %s image reference does not match the transaction", request.TargetContainer) - } - _, err = u.store.CommitBackendContainerDeployment(ctx, transactionID, transaction.BackendContainerDeployment{ - ActivePort: request.TargetPort, - ContainerName: request.TargetContainer, - ImageDigest: request.ImageDigest, - ContainerID: target.ID, - }, "container backend update committed") - if err == nil { - reportProgress(report, Progress{TransactionID: transactionID, State: transaction.StateCommitted, Message: "Container backend update committed"}) - } - return err + return u.drainAndCommitContainer(ctx, transactionID, request, report) case transaction.StateCommitted: return nil case transaction.StateRollingBack: @@ -574,6 +536,53 @@ func (u *Updater) switchAndCommitContainer(ctx context.Context, transactionID st } } +func (u *Updater) switchContainerTraffic(ctx context.Context, transactionID string, request persistedContainerRequest, before hostnginx.Snapshot, after hostnginx.Snapshot, report ProgressReporter) error { + reportProgress(report, Progress{TransactionID: transactionID, State: transaction.StateSwitching, Message: fmt.Sprintf("Switching host Nginx backend traffic to port %d", request.TargetPort)}) + operation := &gatewayOperation{controller: u.gateway, before: before, after: after, receiptPath: request.GatewayReceiptPath} + if _, err := u.coordinator.ExecuteStep(ctx, transactionID, containerGatewaySwitchIntent(request), operation); err != nil { + return u.rollbackContainer(ctx, transactionID, request, before, after, err) + } + _, err := u.store.Transition(ctx, transactionID, transaction.StateVerifying, "host Nginx now routes backend traffic to the healthy container slot") + return err +} + +func (u *Updater) markContainerDraining(ctx context.Context, transactionID string, request persistedContainerRequest) error { + message := "container backend switch verified" + if request.PreviousContainer != "" { + message += "; previous container draining" + } + _, err := u.store.Transition(ctx, transactionID, transaction.StateDraining, message) + return err +} + +func (u *Updater) drainAndCommitContainer(ctx context.Context, transactionID string, request persistedContainerRequest, report ProgressReporter) error { + if request.PreviousContainer != "" { + reportProgress(report, Progress{TransactionID: transactionID, State: transaction.StateDraining, Message: fmt.Sprintf("Draining previous backend container for %s", u.drain)}) + if err := waitContext(ctx, u.drain); err != nil { + return err + } + operation := &containerStopOperation{engine: u.engine, name: request.PreviousContainer, stopGraceSeconds: u.containerStopGraceSeconds} + if _, err := u.coordinator.ExecuteStep(ctx, transactionID, stopPreviousContainerIntent(request), operation); err != nil { + return err + } + } + target, err := u.engine.InspectContainer(ctx, request.TargetContainer) + if err != nil { + return fmt.Errorf("inspect target backend container before commit: %w", err) + } + if !target.Running || target.Dead || target.ID == "" { + return fmt.Errorf("target backend container %s is not running with an exact identity", request.TargetContainer) + } + if target.ImageReference != request.ImmutableReference && target.ImageReference != request.ImageReference { + return fmt.Errorf("target backend container %s image reference does not match the transaction", request.TargetContainer) + } + _, err = u.store.CommitBackendContainerDeployment(ctx, transactionID, transaction.BackendContainerDeployment{ActivePort: request.TargetPort, ContainerName: request.TargetContainer, ImageDigest: request.ImageDigest, ContainerID: target.ID}, "container backend update committed") + if err == nil { + reportProgress(report, Progress{TransactionID: transactionID, State: transaction.StateCommitted, Message: "Container backend update committed"}) + } + return err +} + // rollbackContainer 执行容器后端补偿:恢复 Nginx 上游、停止被补偿的目标容器,并 // 将事务迁入 RolledBack 状态。cause 为触发补偿的原始错误,会与补偿过程中的错误 // 合并后返回。