From bf5c8d138aac0e162aab8d434266dbfda3496d60 Mon Sep 17 00:00:00 2001 From: FogDong Date: Fri, 27 May 2022 17:37:57 +0800 Subject: [PATCH] remove the terminate workflow to pkg and add feature gates Signed-off-by: FogDong --- .../templates/kubevela-controller.yaml | 2 +- .../templates/kubevela-controller.yaml | 2 +- cmd/core/main.go | 1 - pkg/apiserver/domain/service/workflow.go | 40 ++++++++++++++++++- .../application/application_controller.go | 5 ++- .../application_controller_test.go | 10 +++-- pkg/features/controller_features.go | 3 ++ pkg/workflow/tasks/custom/task.go | 8 ++-- pkg/workflow/tasks/discover.go | 6 ++- pkg/workflow/workflow.go | 29 +++++++------- pkg/workflow/workflow_test.go | 11 +++-- references/cli/workflow.go | 32 +-------------- 12 files changed, 83 insertions(+), 66 deletions(-) diff --git a/charts/vela-core/templates/kubevela-controller.yaml b/charts/vela-core/templates/kubevela-controller.yaml index b89a9c20e..bcd4b0d48 100644 --- a/charts/vela-core/templates/kubevela-controller.yaml +++ b/charts/vela-core/templates/kubevela-controller.yaml @@ -169,10 +169,10 @@ spec: - "--concurrent-reconciles={{ .Values.concurrentReconciles }}" - "--kube-api-qps={{ .Values.kubeClient.qps }}" - "--kube-api-burst={{ .Values.kubeClient.burst }}" - - "--enable-suspend-failed-workflow={{ .Values.workflow.enableSuspendFailedWorkflow }}" - "--max-workflow-wait-backoff-time={{ .Values.workflow.backoff.maxTime.waitState }}" - "--max-workflow-failed-backoff-time={{ .Values.workflow.backoff.maxTime.failedState }}" - "--max-workflow-step-error-retry-times={{ .Values.workflow.step.errorRetryTimes }}" + - "--feature-gates=EnableSuspendFailedWorkflow={{- .Values.workflow.enableSuspendFailedWorkflow | toString -}}" - "--feature-gates=AuthenticateApplication={{- .Values.authentication.enabled | toString -}}" {{ if .Values.authentication.enabled }} {{ if .Values.authentication.withUser }} diff --git a/charts/vela-minimal/templates/kubevela-controller.yaml b/charts/vela-minimal/templates/kubevela-controller.yaml index 7d7d705bc..4a59d63da 100644 --- a/charts/vela-minimal/templates/kubevela-controller.yaml +++ b/charts/vela-minimal/templates/kubevela-controller.yaml @@ -139,10 +139,10 @@ spec: - "--concurrent-reconciles={{ .Values.concurrentReconciles }}" - "--kube-api-qps={{ .Values.kubeClient.qps }}" - "--kube-api-burst={{ .Values.kubeClient.burst }}" - - "--enable-suspend-failed-workflow={{ .Values.workflow.enableSuspendFailedWorkflow }}" - "--max-workflow-wait-backoff-time={{ .Values.workflow.backoff.maxTime.waitState }}" - "--max-workflow-failed-backoff-time={{ .Values.workflow.backoff.maxTime.failedState }}" - "--max-workflow-step-error-retry-times={{ .Values.workflow.step.errorRetryTimes }}" + - "--feature-gates=EnableSuspendFailedWorkflow={{- .Values.workflow.enableSuspendFailedWorkflow | toString -}}" - "--feature-gates=AuthenticateApplication={{- .Values.authentication.enabled | toString -}}" {{ if .Values.authentication.enabled }} {{ if .Values.authentication.withUser }} diff --git a/cmd/core/main.go b/cmd/core/main.go index 1120286fc..bee60cc9b 100644 --- a/cmd/core/main.go +++ b/cmd/core/main.go @@ -144,7 +144,6 @@ func main() { flag.IntVar(&workflow.MaxWorkflowWaitBackoffTime, "max-workflow-wait-backoff-time", 60, "Set the max workflow wait backoff time, default is 60") flag.IntVar(&workflow.MaxWorkflowFailedBackoffTime, "max-workflow-failed-backoff-time", 300, "Set the max workflow wait backoff time, default is 300") flag.IntVar(&custom.MaxWorkflowStepErrorRetryTimes, "max-workflow-step-error-retry-times", 10, "Set the max workflow step error retry times, default is 10") - flag.BoolVar(&custom.EnableSuspendFailedWorkflow, "enable-suspend-failed-workflow", false, "Enable suspend failed workflow, defaults to false, if set to true, the if capability in workflow is disabled") utilfeature.DefaultMutableFeatureGate.AddFlag(flag.CommandLine) flag.Parse() diff --git a/pkg/apiserver/domain/service/workflow.go b/pkg/apiserver/domain/service/workflow.go index dcb4e2ce4..944e7f3ab 100644 --- a/pkg/apiserver/domain/service/workflow.go +++ b/pkg/apiserver/domain/service/workflow.go @@ -43,7 +43,7 @@ import ( "github.com/oam-dev/kubevela/pkg/oam/util" utils2 "github.com/oam-dev/kubevela/pkg/utils" "github.com/oam-dev/kubevela/pkg/utils/apply" - "github.com/oam-dev/kubevela/references/cli" + "github.com/oam-dev/kubevela/pkg/workflow/tasks/custom" ) // WorkflowService workflow manage api @@ -578,7 +578,7 @@ func (w *workflowServiceImpl) TerminateRecord(ctx context.Context, appModel *mod return err } - if err := cli.TerminateWorkflow(w.KubeClient, oamApp); err != nil { + if err := TerminateWorkflow(w.KubeClient, oamApp); err != nil { return err } if err := w.syncWorkflowStatus(ctx, oamApp, recordName, oamApp.Name); err != nil { @@ -588,6 +588,42 @@ func (w *workflowServiceImpl) TerminateRecord(ctx context.Context, appModel *mod return nil } +// TerminateWorkflow terminate workflow +func TerminateWorkflow(kubecli client.Client, app *v1beta1.Application) error { + // set the workflow terminated to true + app.Status.Workflow.Terminated = true + steps := app.Status.Workflow.Steps + for i, step := range steps { + switch step.Phase { + case common.WorkflowStepPhaseFailed: + if step.Reason != custom.StatusReasonFailedAfterRetries { + steps[i].Reason = custom.StatusReasonTerminate + } + case common.WorkflowStepPhaseRunning: + steps[i].Phase = common.WorkflowStepPhaseFailed + steps[i].Reason = custom.StatusReasonTerminate + default: + } + for j, sub := range step.SubStepsStatus { + switch sub.Phase { + case common.WorkflowStepPhaseFailed: + if sub.Reason != custom.StatusReasonFailedAfterRetries { + steps[i].SubStepsStatus[j].Phase = custom.StatusReasonTerminate + } + case common.WorkflowStepPhaseRunning: + steps[i].SubStepsStatus[j].Phase = common.WorkflowStepPhaseFailed + steps[i].SubStepsStatus[j].Reason = custom.StatusReasonTerminate + default: + } + } + } + + if err := kubecli.Status().Patch(context.TODO(), app, client.Merge); err != nil { + return err + } + return nil +} + func (w *workflowServiceImpl) RollbackRecord(ctx context.Context, appModel *model.Application, workflow *model.Workflow, recordName, revisionVersion string) error { if revisionVersion == "" { // find the latest complete revision version diff --git a/pkg/controller/core.oam.dev/v1alpha2/application/application_controller.go b/pkg/controller/core.oam.dev/v1alpha2/application/application_controller.go index 5f826f43a..0570be52a 100644 --- a/pkg/controller/core.oam.dev/v1alpha2/application/application_controller.go +++ b/pkg/controller/core.oam.dev/v1alpha2/application/application_controller.go @@ -29,6 +29,7 @@ import ( kerrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apiserver/pkg/util/feature" "k8s.io/client-go/util/workqueue" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" @@ -47,6 +48,7 @@ import ( common2 "github.com/oam-dev/kubevela/pkg/controller/common" core "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev" "github.com/oam-dev/kubevela/pkg/cue/packages" + "github.com/oam-dev/kubevela/pkg/features" monitorContext "github.com/oam-dev/kubevela/pkg/monitor/context" "github.com/oam-dev/kubevela/pkg/monitor/metrics" "github.com/oam-dev/kubevela/pkg/oam" @@ -56,7 +58,6 @@ import ( "github.com/oam-dev/kubevela/pkg/resourcetracker" "github.com/oam-dev/kubevela/pkg/workflow" wfContext "github.com/oam-dev/kubevela/pkg/workflow/context" - "github.com/oam-dev/kubevela/pkg/workflow/tasks/custom" "github.com/oam-dev/kubevela/version" ) @@ -230,7 +231,7 @@ func (r *Reconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Resu handler.app.Status.Workflow.SuspendState = "" return r.gcResourceTrackers(logCtx, handler, common.ApplicationRunningWorkflow, false, false) } - if !workflow.IsFailedAfterRetry(app) || !custom.EnableSuspendFailedWorkflow { + if !workflow.IsFailedAfterRetry(app) || !feature.DefaultMutableFeatureGate.Enabled(features.EnableSuspendFailedWorkflow) { r.stateKeep(logCtx, handler, app) } return r.gcResourceTrackers(logCtx, handler, common.ApplicationWorkflowSuspending, false, true) diff --git a/pkg/controller/core.oam.dev/v1alpha2/application/application_controller_test.go b/pkg/controller/core.oam.dev/v1alpha2/application/application_controller_test.go index 32ced08b0..e7a5490db 100644 --- a/pkg/controller/core.oam.dev/v1alpha2/application/application_controller_test.go +++ b/pkg/controller/core.oam.dev/v1alpha2/application/application_controller_test.go @@ -23,8 +23,11 @@ import ( "net" "net/http" "net/http/httptest" + "testing" "time" + utilfeature "k8s.io/apiserver/pkg/util/feature" + featuregatetesting "k8s.io/component-base/featuregate/testing" "k8s.io/utils/pointer" . "github.com/onsi/ginkgo" @@ -50,6 +53,7 @@ import ( "github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1" stdv1alpha1 "github.com/oam-dev/kubevela/apis/standard.oam.dev/v1alpha1" velatypes "github.com/oam-dev/kubevela/apis/types" + "github.com/oam-dev/kubevela/pkg/features" "github.com/oam-dev/kubevela/pkg/oam" "github.com/oam-dev/kubevela/pkg/oam/testutil" "github.com/oam-dev/kubevela/pkg/oam/util" @@ -1864,7 +1868,7 @@ var _ = Describe("Test Application Controller", func() { }) It("application with dag workflow failed after retries", func() { - custom.EnableSuspendFailedWorkflow = true + defer featuregatetesting.SetFeatureGateDuringTest(&testing.T{}, utilfeature.DefaultFeatureGate, features.EnableSuspendFailedWorkflow, true)() ns := corev1.Namespace{ ObjectMeta: metav1.ObjectMeta{ Name: "dag-failed-after-retries", @@ -1973,7 +1977,7 @@ var _ = Describe("Test Application Controller", func() { }) It("application with step by step workflow failed after retries", func() { - custom.EnableSuspendFailedWorkflow = true + defer featuregatetesting.SetFeatureGateDuringTest(&testing.T{}, utilfeature.DefaultFeatureGate, features.EnableSuspendFailedWorkflow, true)() ns := corev1.Namespace{ ObjectMeta: metav1.ObjectMeta{ Name: "step-by-step-failed-after-retries", @@ -2172,7 +2176,6 @@ var _ = Describe("Test Application Controller", func() { }) It("application with if always in workflow", func() { - custom.EnableSuspendFailedWorkflow = false ns := corev1.Namespace{ ObjectMeta: metav1.ObjectMeta{ Name: "app-with-if-always-workflow", @@ -2272,7 +2275,6 @@ var _ = Describe("Test Application Controller", func() { }) It("application with if always in workflow sub steps", func() { - custom.EnableSuspendFailedWorkflow = false ns := corev1.Namespace{ ObjectMeta: metav1.ObjectMeta{ Name: "app-with-if-always-workflow-sub-steps", diff --git a/pkg/features/controller_features.go b/pkg/features/controller_features.go index eec989cea..5460dbda8 100644 --- a/pkg/features/controller_features.go +++ b/pkg/features/controller_features.go @@ -33,6 +33,8 @@ const ( DeprecatedObjectLabelSelector featuregate.Feature = "DeprecatedObjectLabelSelector" // LegacyResourceTrackerGC enable the gc of legacy resource tracker in managed clusters LegacyResourceTrackerGC featuregate.Feature = "LegacyResourceTrackerGC" + // EnableSuspendFailedWorkflow enable suspend failed workflow + EnableSuspendFailedWorkflow featuregate.Feature = "EnableSuspendFailedWorkflow" // Edge Features @@ -45,6 +47,7 @@ var defaultFeatureGates = map[featuregate.Feature]featuregate.FeatureSpec{ LegacyObjectTypeIdentifier: {Default: false, PreRelease: featuregate.Alpha}, DeprecatedObjectLabelSelector: {Default: false, PreRelease: featuregate.Alpha}, LegacyResourceTrackerGC: {Default: true, PreRelease: featuregate.Alpha}, + EnableSuspendFailedWorkflow: {Default: false, PreRelease: featuregate.Alpha}, AuthenticateApplication: {Default: false, PreRelease: featuregate.Alpha}, } diff --git a/pkg/workflow/tasks/custom/task.go b/pkg/workflow/tasks/custom/task.go index dc5621984..7c4151914 100644 --- a/pkg/workflow/tasks/custom/task.go +++ b/pkg/workflow/tasks/custom/task.go @@ -24,6 +24,7 @@ import ( "cuelang.org/go/cue" "github.com/pkg/errors" + "k8s.io/apiserver/pkg/util/feature" "github.com/oam-dev/kubevela/apis/core.oam.dev/common" "github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1" @@ -32,6 +33,7 @@ import ( "github.com/oam-dev/kubevela/pkg/cue/model/value" "github.com/oam-dev/kubevela/pkg/cue/packages" "github.com/oam-dev/kubevela/pkg/cue/process" + "github.com/oam-dev/kubevela/pkg/features" monitorContext "github.com/oam-dev/kubevela/pkg/monitor/context" wfContext "github.com/oam-dev/kubevela/pkg/workflow/context" "github.com/oam-dev/kubevela/pkg/workflow/hooks" @@ -42,8 +44,6 @@ import ( var ( // MaxWorkflowStepErrorRetryTimes is the max retry times of the failed workflow step. MaxWorkflowStepErrorRetryTimes = 10 - // EnableSuspendFailedWorkflow enable suspend failed workflow - EnableSuspendFailedWorkflow = false ) const ( @@ -157,7 +157,7 @@ func (t *TaskLoader) makeTaskGenerator(templ string) (wfTypes.TaskGenerator, err return CheckPending(ctx, wfStep, stepStatus) } tRunner.skip = func(dependsOnPhase common.WorkflowStepPhase, stepStatus map[string]common.StepStatus) (common.StepStatus, bool) { - if EnableSuspendFailedWorkflow { + if feature.DefaultMutableFeatureGate.Enabled(features.EnableSuspendFailedWorkflow) { return exec.status(), false } skip := SkipTaskRunner(&SkipOptions{ @@ -506,7 +506,7 @@ func CheckPending(ctx wfContext.Context, step v1beta1.WorkflowStep, stepStatus m // IsStepFinish will decide whether step is finish. func IsStepFinish(phase common.WorkflowStepPhase, reason string) bool { - if EnableSuspendFailedWorkflow { + if feature.DefaultMutableFeatureGate.Enabled(features.EnableSuspendFailedWorkflow) { return phase == common.WorkflowStepPhaseSucceeded } if phase == common.WorkflowStepPhaseFailed { diff --git a/pkg/workflow/tasks/discover.go b/pkg/workflow/tasks/discover.go index 791f715eb..38efb6743 100644 --- a/pkg/workflow/tasks/discover.go +++ b/pkg/workflow/tasks/discover.go @@ -22,6 +22,7 @@ import ( builtintime "time" "github.com/pkg/errors" + "k8s.io/apiserver/pkg/util/feature" "k8s.io/client-go/rest" "sigs.k8s.io/controller-runtime/pkg/client" @@ -29,6 +30,7 @@ import ( "github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1" "github.com/oam-dev/kubevela/pkg/cue/packages" "github.com/oam-dev/kubevela/pkg/cue/process" + "github.com/oam-dev/kubevela/pkg/features" monitorContext "github.com/oam-dev/kubevela/pkg/monitor/context" "github.com/oam-dev/kubevela/pkg/oam/discoverymapper" "github.com/oam-dev/kubevela/pkg/velaql/providers/query" @@ -160,7 +162,7 @@ func (tr *suspendTaskRunner) Skip(dependsOnPhase common.WorkflowStepPhase, stepS Type: types.WorkflowStepTypeSuspend, Phase: tr.phase, } - if custom.EnableSuspendFailedWorkflow { + if feature.DefaultMutableFeatureGate.Enabled(features.EnableSuspendFailedWorkflow) { return status, false } skip := custom.SkipTaskRunner(&custom.SkipOptions{ @@ -197,7 +199,7 @@ func (tr *stepGroupTaskRunner) Skip(dependsOnPhase common.WorkflowStepPhase, ste Name: tr.step.Name, Type: types.WorkflowStepTypeStepGroup, } - if custom.EnableSuspendFailedWorkflow { + if feature.DefaultMutableFeatureGate.Enabled(features.EnableSuspendFailedWorkflow) { return status, false } skip := custom.SkipTaskRunner(&custom.SkipOptions{ diff --git a/pkg/workflow/workflow.go b/pkg/workflow/workflow.go index e86328b58..6482dc7d8 100644 --- a/pkg/workflow/workflow.go +++ b/pkg/workflow/workflow.go @@ -25,6 +25,7 @@ import ( "github.com/pkg/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apiserver/pkg/util/feature" "sigs.k8s.io/controller-runtime/pkg/client" "github.com/oam-dev/kubevela/apis/core.oam.dev/common" @@ -32,6 +33,7 @@ import ( oamcore "github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1" "github.com/oam-dev/kubevela/pkg/controller/utils" "github.com/oam-dev/kubevela/pkg/cue/model/value" + "github.com/oam-dev/kubevela/pkg/features" monitorContext "github.com/oam-dev/kubevela/pkg/monitor/context" "github.com/oam-dev/kubevela/pkg/monitor/metrics" "github.com/oam-dev/kubevela/pkg/oam" @@ -153,6 +155,7 @@ func (w *workflow) ExecuteSteps(ctx monitorContext.Context, appRev *oamcore.Appl } e.checkWorkflowStatusMessage(wfStatus) + fmt.Println(99999, e.status.Message) StepStatusCache.Store(cacheKey, len(wfStatus.Steps)) allTasksDone, allTasksSucceeded = w.allDone(taskRunners) if wfStatus.Terminated { @@ -536,20 +539,16 @@ func (e *engine) Run(taskRunners []wfTypes.TaskRunner, dag bool) error { } func (e *engine) checkWorkflowStatusMessage(wfStatus *common.WorkflowStatus) { - if !e.waiting && e.failedAfterRetries { - if custom.EnableSuspendFailedWorkflow { - e.status.Message = MessageSuspendFailedAfterRetries - } else { - e.status.Message = MessageTerminatedFailedAfterRetries - } - return - } - - if wfStatus.Terminated { + switch { + case !e.waiting && e.failedAfterRetries && feature.DefaultMutableFeatureGate.Enabled(features.EnableSuspendFailedWorkflow): + e.status.Message = MessageSuspendFailedAfterRetries + case e.failedAfterRetries && !feature.DefaultMutableFeatureGate.Enabled(features.EnableSuspendFailedWorkflow): + e.status.Message = MessageTerminatedFailedAfterRetries + case wfStatus.Terminated: e.status.Message = string(common.WorkflowStateTerminated) - } - if wfStatus.Suspend { + case wfStatus.Suspend: e.status.Message = string(common.WorkflowStateSuspended) + default: } } @@ -696,16 +695,16 @@ func (e *engine) updateStepStatus(status common.StepStatus) { } func (e *engine) checkFailedAfterRetries() { - if !e.waiting && e.failedAfterRetries && custom.EnableSuspendFailedWorkflow { + if !e.waiting && e.failedAfterRetries && feature.DefaultMutableFeatureGate.Enabled(features.EnableSuspendFailedWorkflow) { e.status.Suspend = true } - if e.failedAfterRetries && !custom.EnableSuspendFailedWorkflow { + if e.failedAfterRetries && !feature.DefaultMutableFeatureGate.Enabled(features.EnableSuspendFailedWorkflow) { e.status.Terminated = true } } func (e *engine) needStop() bool { - if custom.EnableSuspendFailedWorkflow { + if feature.DefaultMutableFeatureGate.Enabled(features.EnableSuspendFailedWorkflow) { e.checkFailedAfterRetries() } // if the workflow is terminated, we still need to execute all the remaining steps diff --git a/pkg/workflow/workflow_test.go b/pkg/workflow/workflow_test.go index 3dfe44368..21f40c940 100644 --- a/pkg/workflow/workflow_test.go +++ b/pkg/workflow/workflow_test.go @@ -20,6 +20,7 @@ import ( "context" "encoding/json" "math" + "testing" . "github.com/onsi/ginkgo" . "github.com/onsi/gomega" @@ -29,11 +30,14 @@ import ( corev1 "k8s.io/api/core/v1" kerrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + utilfeature "k8s.io/apiserver/pkg/util/feature" + featuregatetesting "k8s.io/component-base/featuregate/testing" "sigs.k8s.io/yaml" "github.com/oam-dev/kubevela/apis/core.oam.dev/common" oamcore "github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1" "github.com/oam-dev/kubevela/pkg/cue/model/value" + "github.com/oam-dev/kubevela/pkg/features" monitorContext "github.com/oam-dev/kubevela/pkg/monitor/context" wfContext "github.com/oam-dev/kubevela/pkg/workflow/context" "github.com/oam-dev/kubevela/pkg/workflow/tasks" @@ -519,7 +523,7 @@ var _ = Describe("Test Workflow", func() { It("Workflow test for failed after retries with suspend", func() { By("Test failed-after-retries in StepByStep mode with suspend") - custom.EnableSuspendFailedWorkflow = true + defer featuregatetesting.SetFeatureGateDuringTest(&testing.T{}, utilfeature.DefaultFeatureGate, features.EnableSuspendFailedWorkflow, true)() app, runners := makeTestCase([]oamcore.WorkflowStep{ { Name: "s1", @@ -626,7 +630,6 @@ var _ = Describe("Test Workflow", func() { It("Workflow test if always", func() { By("Test if always in StepByStep mode") - custom.EnableSuspendFailedWorkflow = false app, runners := makeTestCase([]oamcore.WorkflowStep{ { Name: "s1", @@ -801,7 +804,7 @@ var _ = Describe("Test Workflow", func() { It("Test failed after retries with sub steps", func() { By("Test failed-after-retries with step group in StepByStep mode") - custom.EnableSuspendFailedWorkflow = true + defer featuregatetesting.SetFeatureGateDuringTest(&testing.T{}, utilfeature.DefaultFeatureGate, features.EnableSuspendFailedWorkflow, true)() app, runners := makeTestCase([]oamcore.WorkflowStep{ { Name: "s1", @@ -1396,7 +1399,7 @@ func makeRunner(name, tpy, ifDecl string, dependsOn []string, subTaskRunners []w Name: name, Type: tpy, } - if custom.EnableSuspendFailedWorkflow { + if utilfeature.DefaultMutableFeatureGate.Enabled(features.EnableSuspendFailedWorkflow) { return status, false } skip := custom.SkipTaskRunner(&custom.SkipOptions{ diff --git a/references/cli/workflow.go b/references/cli/workflow.go index f042f7ed6..aca0b9021 100644 --- a/references/cli/workflow.go +++ b/references/cli/workflow.go @@ -29,6 +29,7 @@ import ( oamcommon "github.com/oam-dev/kubevela/apis/core.oam.dev/common" "github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1" "github.com/oam-dev/kubevela/apis/types" + "github.com/oam-dev/kubevela/pkg/apiserver/domain/service" "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/application" "github.com/oam-dev/kubevela/pkg/controller/utils" "github.com/oam-dev/kubevela/pkg/oam" @@ -36,7 +37,6 @@ import ( "github.com/oam-dev/kubevela/pkg/utils/common" velaerrors "github.com/oam-dev/kubevela/pkg/utils/errors" cmdutil "github.com/oam-dev/kubevela/pkg/utils/util" - "github.com/oam-dev/kubevela/pkg/workflow/tasks/custom" "github.com/oam-dev/kubevela/references/appfile" ) @@ -289,35 +289,7 @@ func resumeWorkflow(kubecli client.Client, app *v1beta1.Application) error { // TerminateWorkflow terminate workflow func TerminateWorkflow(kubecli client.Client, app *v1beta1.Application) error { - // set the workflow terminated to true - app.Status.Workflow.Terminated = true - steps := app.Status.Workflow.Steps - for i, step := range steps { - switch step.Phase { - case oamcommon.WorkflowStepPhaseFailed: - if step.Reason != custom.StatusReasonFailedAfterRetries { - steps[i].Reason = custom.StatusReasonTerminate - } - case oamcommon.WorkflowStepPhaseRunning: - steps[i].Phase = oamcommon.WorkflowStepPhaseFailed - steps[i].Reason = custom.StatusReasonTerminate - default: - } - for j, sub := range step.SubStepsStatus { - switch sub.Phase { - case oamcommon.WorkflowStepPhaseFailed: - if sub.Reason != custom.StatusReasonFailedAfterRetries { - steps[i].SubStepsStatus[j].Phase = custom.StatusReasonTerminate - } - case oamcommon.WorkflowStepPhaseRunning: - steps[i].SubStepsStatus[j].Phase = oamcommon.WorkflowStepPhaseFailed - steps[i].SubStepsStatus[j].Reason = custom.StatusReasonTerminate - default: - } - } - } - - if err := kubecli.Status().Patch(context.TODO(), app, client.Merge); err != nil { + if err := service.TerminateWorkflow(kubecli, app); err != nil { return err }