remove the terminate workflow to pkg and add feature gates

Signed-off-by: FogDong <dongtianxin.tx@alibaba-inc.com>
This commit is contained in:
FogDong
2022-05-27 17:37:57 +08:00
parent d1dd602cc8
commit bf5c8d138a
12 changed files with 83 additions and 66 deletions
@@ -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 }}
@@ -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 }}
-1
View File
@@ -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()
+38 -2
View File
@@ -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
@@ -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)
@@ -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",
+3
View File
@@ -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},
}
+4 -4
View File
@@ -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 {
+4 -2
View File
@@ -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{
+14 -15
View File
@@ -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
+7 -4
View File
@@ -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{
+2 -30
View File
@@ -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
}