Workflow Support Resource GC (#1970)

* gc

* add test cases

* test case
This commit is contained in:
Jian.Li
2021-07-27 19:22:05 +08:00
committed by GitHub
parent 804024599b
commit a736b1f7b0
12 changed files with 237 additions and 77 deletions
+9 -4
View File
@@ -23,6 +23,9 @@ import (
"regexp"
"strings"
"github.com/oam-dev/kubevela/pkg/workflow/providers"
"github.com/oam-dev/kubevela/pkg/workflow/providers/kube"
"github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1"
"github.com/oam-dev/kubevela/pkg/cue/packages"
"github.com/oam-dev/kubevela/pkg/oam/discoverymapper"
@@ -192,13 +195,13 @@ type Appfile struct {
}
// GenerateWorkflowAndPolicy generates workflow steps and policies from an appFile
func (af *Appfile) GenerateWorkflowAndPolicy(ctx context.Context, m discoverymapper.DiscoveryMapper, cli client.Client, pd *packages.PackageDiscover) (policies []*unstructured.Unstructured, steps []wfTypes.TaskRunner, err error) {
func (af *Appfile) GenerateWorkflowAndPolicy(ctx context.Context, m discoverymapper.DiscoveryMapper, cli client.Client, pd *packages.PackageDiscover, dispatcher kube.Dispatcher) (policies []*unstructured.Unstructured, steps []wfTypes.TaskRunner, err error) {
policies, err = af.generateUnstructureds(af.Policies)
if err != nil {
return
}
steps, err = af.generateSteps(ctx, m, cli, pd)
steps, err = af.generateSteps(ctx, m, cli, pd, dispatcher)
return
}
@@ -214,7 +217,7 @@ func (af *Appfile) generateUnstructureds(workloads []*Workload) ([]*unstructured
return uns, nil
}
func (af *Appfile) generateSteps(ctx context.Context, dm discoverymapper.DiscoveryMapper, cli client.Client, pd *packages.PackageDiscover) ([]wfTypes.TaskRunner, error) {
func (af *Appfile) generateSteps(ctx context.Context, dm discoverymapper.DiscoveryMapper, cli client.Client, pd *packages.PackageDiscover, dispatcher kube.Dispatcher) ([]wfTypes.TaskRunner, error) {
loadTaskTemplate := func(ctx context.Context, name string) (string, error) {
templ, err := LoadTemplate(ctx, dm, cli, name, types.TypeWorkflowStep)
if err != nil {
@@ -227,7 +230,9 @@ func (af *Appfile) generateSteps(ctx context.Context, dm discoverymapper.Discove
return "", errors.New("custom workflowStep only support cue")
}
taskDiscover := tasks.NewTaskDiscover(cli, pd, loadTaskTemplate)
handlerProviders := providers.NewProviders()
kube.Install(handlerProviders, cli, dispatcher)
taskDiscover := tasks.NewTaskDiscover(handlerProviders, pd, loadTaskTemplate)
var tasks []wfTypes.TaskRunner
for _, step := range af.WorkflowSteps {
genTask, err := taskDiscover.GetTaskGenerator(ctx, step.Type)
+3 -3
View File
@@ -390,7 +390,7 @@ wait: op.#ConditionalWait & {
},
}
ctx := context.WithValue(context.Background(), util.AppDefinitionNamespace, "default")
runners, err := appfile.generateSteps(ctx, dm, k8sClient, pd)
runners, err := appfile.generateSteps(ctx, dm, k8sClient, pd, nil)
Expect(err).To(BeNil())
Expect(len(runners)).Should(BeEquivalentTo(1))
@@ -404,7 +404,7 @@ wait: op.#ConditionalWait & {
Type: "empty",
},
}
_, err = appfile.generateSteps(ctx, dm, k8sClient, pd)
_, err = appfile.generateSteps(ctx, dm, k8sClient, pd, nil)
Expect(err).NotTo(BeNil())
appfile.WorkflowSteps = []v1beta1.WorkflowStep{
@@ -413,7 +413,7 @@ wait: op.#ConditionalWait & {
Type: "not-cue",
},
}
_, err = appfile.generateSteps(ctx, dm, k8sClient, pd)
_, err = appfile.generateSteps(ctx, dm, k8sClient, pd, nil)
Expect(err).NotTo(BeNil())
})
})
@@ -150,7 +150,7 @@ func (r *Reconciler) Reconcile(req ctrl.Request) (ctrl.Result, error) {
r.Recorder.Event(app, event.Normal(velatypes.ReasonRevisoned, velatypes.MessageRevisioned))
klog.Info("Successfully apply application revision", "application", klog.KObj(app))
policies, wfSteps, err := appFile.GenerateWorkflowAndPolicy(ctx, r.dm, r.Client, r.pd)
policies, wfSteps, err := appFile.GenerateWorkflowAndPolicy(ctx, r.dm, r.Client, r.pd, handler.Dispatch)
if err != nil {
klog.Error(err, "[Handle GenerateWorkflowAndPolicy]")
r.Recorder.Event(app, event.Warning(velatypes.ReasonFailedRender, err))
@@ -174,16 +174,37 @@ func (r *Reconciler) Reconcile(req ctrl.Request) (ctrl.Result, error) {
r.Recorder.Event(app, event.Normal(velatypes.ReasonApplied, velatypes.MessageApplied))
klog.Info("Successfully apply application manifests", "application", klog.KObj(app))
done, err := workflow.NewWorkflow(app, r.Client).ExecuteSteps(ctx, handler.currentAppRev.Name, wfSteps)
done, pause, err := workflow.NewWorkflow(app, r.Client).ExecuteSteps(ctx, handler.currentAppRev.Name, wfSteps)
if err != nil {
klog.Error(err, "[handle workflow]")
r.Recorder.Event(app, event.Warning(velatypes.ReasonFailedWorkflow, err))
return r.endWithNegativeCondition(ctx, app, utils.ErrorCondition("Workflow", err))
}
if pause {
if err := r.patchStatus(ctx, app); err != nil {
return r.endWithNegativeCondition(ctx, app, v1alpha1.ReconcileError(err))
}
return ctrl.Result{}, nil
}
if !done {
return reconcile.Result{RequeueAfter: WorkflowReconcileWaitTime}, r.patchStatus(ctx, app)
}
if wfStatus := app.Status.Workflow; wfStatus != nil && !wfStatus.Terminated {
ref, err := handler.DispatchAndGC(ctx)
if err != nil {
klog.ErrorS(err, "Failed to gc after workflow",
"application", klog.KObj(app))
r.Recorder.Event(app, event.Warning(velatypes.ReasonFailedGC, err))
return r.endWithNegativeCondition(ctx, app, utils.ErrorCondition("GCAfterWorkflow", err))
}
wfStatus.Terminated = true
app.Status.ResourceTracker = ref
return r.endWithNegativeCondition(ctx, app, utils.ReadyCondition("GCAfterWorkflow"))
}
// if inplace is false and rolloutPlan is nil, it means the user will use an outer AppRollout object to rollout the application
if handler.app.Spec.RolloutPlan != nil {
res, err := handler.handleRollout(ctx)
@@ -20,7 +20,6 @@ import (
"context"
runtimev1alpha1 "github.com/crossplane/crossplane-runtime/apis/core/v1alpha1"
"github.com/pkg/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
@@ -49,10 +48,51 @@ type AppHandler struct {
app *v1beta1.Application
currentAppRev *v1beta1.ApplicationRevision
latestAppRev *v1beta1.ApplicationRevision
latestTracker *v1beta1.ResourceTracker
dispatcher *dispatch.AppManifestsDispatcher
isNewRevision bool
currentRevHash string
}
// Dispatch apply manifests into k8s.
func (h *AppHandler) Dispatch(ctx context.Context, manifests ...*unstructured.Unstructured) error {
h.initDispatcher()
_, err := h.dispatcher.Dispatch(ctx, manifests)
return err
}
// DispatchAndGC apply manifests and do GC.
func (h *AppHandler) DispatchAndGC(ctx context.Context, manifests ...*unstructured.Unstructured) (*runtimev1alpha1.TypedReference, error) {
h.initDispatcher()
tracker, err := h.dispatcher.EndAndGC(h.latestTracker).Dispatch(ctx, manifests)
if err != nil {
return nil, errors.WithMessage(err, "cannot dispatch application manifests")
}
return &runtimev1alpha1.TypedReference{
APIVersion: tracker.APIVersion,
Kind: tracker.Kind,
Name: tracker.Name,
UID: tracker.UID,
}, nil
}
func (h *AppHandler) initDispatcher() {
if h.latestTracker == nil {
if h.app.Status.ResourceTracker != nil {
h.latestTracker = &v1beta1.ResourceTracker{}
h.latestTracker.Name = h.app.Status.ResourceTracker.Name
} else if h.app.Status.LatestRevision != nil {
h.latestTracker = &v1beta1.ResourceTracker{}
h.latestTracker.SetName(dispatch.ConstructResourceTrackerName(h.app.Status.LatestRevision.Name, h.app.Namespace))
}
}
if h.dispatcher == nil {
// only do GC when ALL resources are dispatched successfully
// so skip GC while dispatching addon resources
h.dispatcher = dispatch.NewAppManifestsDispatcher(h.r.Client, h.currentAppRev).StartAndSkipGC(h.latestTracker)
}
}
// ApplyAppManifests will dispatch Application manifests
func (h *AppHandler) ApplyAppManifests(ctx context.Context, comps []*types.ComponentManifest, policies []*unstructured.Unstructured) error {
appRev := h.currentAppRev
@@ -68,13 +108,10 @@ func (h *AppHandler) ApplyAppManifests(ctx context.Context, comps []*types.Compo
latestTracker = &v1beta1.ResourceTracker{}
latestTracker.SetName(dispatch.ConstructResourceTrackerName(h.app.Status.LatestRevision.Name, h.app.Namespace))
}
// only do GC when ALL resources are dispatched successfully
// so skip GC while dispatching addon resources
d := dispatch.NewAppManifestsDispatcher(h.r.Client, appRev).StartAndSkipGC(latestTracker)
// dispatch packaged workload resources before dispatching assembled manifests
for _, comp := range comps {
if len(comp.PackagedWorkloadResources) != 0 {
if _, err := d.Dispatch(ctx, comp.PackagedWorkloadResources); err != nil {
if err := h.Dispatch(ctx, comp.PackagedWorkloadResources...); err != nil {
return errors.WithMessage(err, "cannot dispatch packaged workload resources")
}
}
@@ -82,15 +119,13 @@ func (h *AppHandler) ApplyAppManifests(ctx context.Context, comps []*types.Compo
continue
}
}
a := assemble.NewAppManifests(appRev).WithWorkloadOption(assemble.DiscoveryHelmBasedWorkload(ctx, h.r.Client))
a := assemble.NewAppManifests(h.currentAppRev).WithWorkloadOption(assemble.DiscoveryHelmBasedWorkload(ctx, h.r.Client))
manifests, err := a.AssembledManifests()
if err != nil {
return errors.WithMessage(err, "cannot assemble application manifests")
}
if _, err := d.EndAndGC(latestTracker).Dispatch(ctx, manifests); err != nil {
return errors.WithMessage(err, "cannot dispatch application manifests")
}
return nil
_, err = h.DispatchAndGC(ctx, manifests...)
return err
}
func (h *AppHandler) aggregateHealthStatus(appFile *appfile.Appfile) ([]common.ApplicationComponentStatus, bool, error) {
@@ -111,7 +111,7 @@ var _ = Describe("Test Workflow", func() {
})
It("should execute workflow step to apply and wait", func() {
Expect(k8sClient.Create(ctx, appWithWorkflow)).Should(BeNil())
Expect(k8sClient.Create(ctx, appWithWorkflow.DeepCopy())).Should(BeNil())
// first try to add finalizer
tryReconcile(reconciler, appWithWorkflow.Name, appWithWorkflow.Namespace)
@@ -146,6 +146,7 @@ var _ = Describe("Test Workflow", func() {
triggerWorkflowStepToSucceed(stepObj)
Expect(k8sClient.Update(ctx, stepObj)).Should(BeNil())
tryReconcile(reconciler, appWithWorkflow.Name, appWithWorkflow.Namespace)
tryReconcile(reconciler, appWithWorkflow.Name, appWithWorkflow.Namespace)
// check workflow status is succeeded
@@ -155,7 +156,50 @@ var _ = Describe("Test Workflow", func() {
}, appObj)).Should(BeNil())
Expect(appObj.Status.Workflow.Steps[0].Phase).Should(Equal(common.WorkflowStepPhaseSucceeded))
Expect(appObj.Status.Workflow.Terminated).Should(BeTrue())
})
It("test workflow suspend", func() {
suspendApp := appWithWorkflow.DeepCopy()
suspendApp.Name = "test-app-suspend"
suspendApp.Spec.Workflow.Steps = []oamcore.WorkflowStep{{
Name: "suspend",
Type: "suspend",
Properties: runtime.RawExtension{Raw: []byte(`{}`)},
}}
Expect(k8sClient.Create(ctx, suspendApp)).Should(BeNil())
// first try to add finalizer
tryReconcile(reconciler, suspendApp.Name, suspendApp.Namespace)
tryReconcile(reconciler, suspendApp.Name, suspendApp.Namespace)
appObj := &oamcore.Application{}
Expect(k8sClient.Get(ctx, client.ObjectKey{
Name: suspendApp.Name,
Namespace: suspendApp.Namespace,
}, appObj)).Should(BeNil())
Expect(appObj.Status.Workflow.Suspend).Should(BeTrue())
Expect(appObj.Status.Phase).Should(BeEquivalentTo(common.ApplicationRunningWorkflow))
// resume
appObj.Status.Workflow.Suspend = false
Expect(k8sClient.Status().Patch(ctx, appObj, client.Merge)).Should(BeNil())
tryReconcile(reconciler, suspendApp.Name, suspendApp.Namespace)
tryReconcile(reconciler, suspendApp.Name, suspendApp.Namespace)
appObj = &oamcore.Application{}
Expect(k8sClient.Get(ctx, client.ObjectKey{
Name: suspendApp.Name,
Namespace: suspendApp.Namespace,
}, appObj)).Should(BeNil())
Expect(appObj.Status.Workflow.Suspend).Should(BeFalse())
Expect(appObj.Status.Workflow.Terminated).Should(BeTrue())
Expect(appObj.Status.Workflow.StepIndex).Should(BeEquivalentTo(1))
})
})
func triggerWorkflowStepToSucceed(obj *unstructured.Unstructured) {
+1 -1
View File
@@ -25,5 +25,5 @@ import (
type Workflow interface {
// ExecuteSteps executes the steps of an Application with given steps of rendered resources.
// It returns done=true only if all steps are executed and succeeded.
ExecuteSteps(ctx context.Context, appRevName string, taskRunners []types.TaskRunner) (done bool, err error)
ExecuteSteps(ctx context.Context, appRevName string, taskRunners []types.TaskRunner) (done bool, pause bool, err error)
}
+10 -9
View File
@@ -19,12 +19,10 @@ package kube
import (
"context"
"github.com/oam-dev/kubevela/pkg/cue/model/value"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"sigs.k8s.io/controller-runtime/pkg/client"
"github.com/oam-dev/kubevela/pkg/utils/apply"
"github.com/oam-dev/kubevela/pkg/cue/model/value"
wfContext "github.com/oam-dev/kubevela/pkg/workflow/context"
"github.com/oam-dev/kubevela/pkg/workflow/providers"
"github.com/oam-dev/kubevela/pkg/workflow/types"
@@ -35,9 +33,12 @@ const (
ProviderName = "kube"
)
// Dispatcher is a client for apply resources.
type Dispatcher func(ctx context.Context, manifests ...*unstructured.Unstructured) error
type provider struct {
deploy *apply.APIApplicator
cli client.Client
apply func(ctx context.Context, manifests ...*unstructured.Unstructured) error
cli client.Client
}
// Apply create or update CR in cluster.
@@ -51,7 +52,7 @@ func (h *provider) Apply(ctx wfContext.Context, v *value.Value, act types.Action
if workload.GetNamespace() == "" {
workload.SetNamespace("default")
}
if err := h.deploy.Apply(deployCtx, workload); err != nil {
if err := h.apply(deployCtx, workload); err != nil {
return err
}
return v.FillObject(workload.Object)
@@ -77,10 +78,10 @@ func (h *provider) Read(ctx wfContext.Context, v *value.Value, act types.Action)
}
// Install register handlers to provider discover.
func Install(p providers.Providers, cli client.Client) {
func Install(p providers.Providers, cli client.Client, apply Dispatcher) {
prd := &provider{
deploy: apply.NewAPIApplicator(cli),
cli: cli,
apply: apply,
cli: cli,
}
p.Register(ProviderName, map[string]providers.Handler{
"apply": prd.Apply,
+14 -6
View File
@@ -22,6 +22,8 @@ import (
"testing"
"time"
"k8s.io/apimachinery/pkg/api/errors"
. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
@@ -41,7 +43,6 @@ import (
"github.com/oam-dev/kubevela/pkg/cue/model/value"
"github.com/oam-dev/kubevela/pkg/cue/packages"
"github.com/oam-dev/kubevela/pkg/utils/apply"
wfContext "github.com/oam-dev/kubevela/pkg/workflow/context"
)
@@ -57,8 +58,18 @@ var pd *packages.PackageDiscover
var _ = Describe("Test Workflow Provider Kube", func() {
It("apply and read", func() {
p := &provider{
deploy: apply.NewAPIApplicator(k8sClient),
cli: k8sClient,
apply: func(ctx context.Context, manifests ...*unstructured.Unstructured) error {
for _, obj := range manifests {
if err := k8sClient.Create(ctx, obj); err != nil {
if errors.IsAlreadyExists(err) {
return k8sClient.Update(ctx, obj)
}
return err
}
}
return nil
},
cli: k8sClient,
}
ctx, err := newWorkflowContextForTest()
Expect(err).ToNot(HaveOccurred())
@@ -138,9 +149,6 @@ metadata: {
labels: {
app: "nginx"
}
annotations: {
"app.oam.dev/last-applied-configuration": "{\"apiVersion\":\"v1\",\"kind\":\"Pod\",\"metadata\":{\"labels\":{\"app\":\"nginx\"},\"name\":\"app\",\"namespace\":\"default\"},\"spec\":{\"containers\":[{\"env\":[{\"name\":\"APP\",\"value\":\"nginx\"}],\"image\":\"nginx:1.14.2\",\"imagePullPolicy\":\"IfNotPresent\",\"name\":\"main\",\"ports\":[{\"containerPort\":8080,\"protocol\":\"TCP\"}]}]}}"
}
namespace: "default"
resourceVersion: "44"
selfLink: "/api/v1/namespaces/default/pods/app"
+1 -1
View File
@@ -40,7 +40,7 @@ var k8sClient client.Client
var testEnv *envtest.Environment
var scheme = runtime.NewScheme()
func TestDefinition(t *testing.T) {
func TestWorkflow(t *testing.T) {
RegisterFailHandler(Fail)
RunSpecsWithDefaultAndCustomReporters(t,
+2 -6
View File
@@ -20,14 +20,12 @@ import (
"context"
"github.com/pkg/errors"
"sigs.k8s.io/controller-runtime/pkg/client"
"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/pkg/cue/packages"
wfContext "github.com/oam-dev/kubevela/pkg/workflow/context"
"github.com/oam-dev/kubevela/pkg/workflow/providers"
"github.com/oam-dev/kubevela/pkg/workflow/providers/kube"
"github.com/oam-dev/kubevela/pkg/workflow/providers/workspace"
"github.com/oam-dev/kubevela/pkg/workflow/tasks/custom"
"github.com/oam-dev/kubevela/pkg/workflow/types"
@@ -68,10 +66,8 @@ func suspend(step v1beta1.WorkflowStep) (types.TaskRunner, error) {
}
// NewTaskDiscover will create a client for load task generator.
func NewTaskDiscover(cli client.Client, pd *packages.PackageDiscover, loadTemplate custom.LoadTaskTemplate) types.TaskDiscover {
providerHandlers := providers.NewProviders()
kube.Install(providerHandlers, cli)
func NewTaskDiscover(providerHandlers providers.Providers, pd *packages.PackageDiscover, loadTemplate custom.LoadTaskTemplate) types.TaskDiscover {
// install builtin provider
workspace.Install(providerHandlers)
return &taskDiscover{
+44 -25
View File
@@ -42,14 +42,14 @@ func NewWorkflow(app *oamcore.Application, cli client.Client) Workflow {
}
// ExecuteSteps process workflow step in order.
func (w *workflow) ExecuteSteps(ctx context.Context, rev string, taskRunners []wfTypes.TaskRunner) (bool, error) {
func (w *workflow) ExecuteSteps(ctx context.Context, rev string, taskRunners []wfTypes.TaskRunner) (done bool, pause bool, gerr error) {
if w.app.Spec.Workflow == nil {
return true, nil
return true, false, nil
}
steps := w.app.Spec.Workflow.Steps
if len(steps) == 0 {
return true, nil
return true, false, nil
}
if w.app.Status.Workflow == nil || w.app.Status.Workflow.AppRevision != rev {
@@ -61,25 +61,38 @@ func (w *workflow) ExecuteSteps(ctx context.Context, rev string, taskRunners []w
wfStatus := w.app.Status.Workflow
if len(taskRunners) <= wfStatus.StepIndex || wfStatus.Terminated || wfStatus.Suspend {
return true, nil
if wfStatus.Terminated {
done = true
return
}
w.app.Status.Phase = common.ApplicationRunningWorkflow
if len(taskRunners) <= wfStatus.StepIndex {
done = true
return
}
if wfStatus.Suspend {
pause = true
return
}
var (
wfCtx wfContext.Context
err error
)
if wfStatus.ContextBackend != nil {
wfCtx, err = wfContext.LoadContext(w.cli, w.app.Namespace, rev)
if err != nil {
return false, errors.WithMessage(err, "load context")
wfCtx, gerr = wfContext.LoadContext(w.cli, w.app.Namespace, rev)
if gerr != nil {
gerr = errors.WithMessage(gerr, "load context")
return
}
} else {
wfCtx, err = wfContext.NewContext(w.cli, w.app.Namespace, rev)
if err != nil {
return false, errors.WithMessage(err, "new context")
wfCtx, gerr = wfContext.NewContext(w.cli, w.app.Namespace, rev)
if gerr != nil {
gerr = errors.WithMessage(gerr, "new context")
return
}
wfStatus.ContextBackend = wfCtx.StoreRef()
}
@@ -87,7 +100,8 @@ func (w *workflow) ExecuteSteps(ctx context.Context, rev string, taskRunners []w
for _, run := range taskRunners[wfStatus.StepIndex:] {
status, operation, err := run(wfCtx)
if err != nil {
return false, err
gerr = err
return
}
var conditionUpdated bool
@@ -103,24 +117,29 @@ func (w *workflow) ExecuteSteps(ctx context.Context, rev string, taskRunners []w
wfStatus.Steps = append(wfStatus.Steps, status)
}
if status.Phase != common.WorkflowStepPhaseSucceeded {
return
}
if err := wfCtx.Commit(); err != nil {
gerr = errors.WithMessage(err, "commit workflow context")
return
}
wfStatus.StepIndex++
if operation != nil {
wfStatus.Terminated = operation.Terminated
wfStatus.Suspend = operation.Suspend
}
if status.Phase != common.WorkflowStepPhaseSucceeded {
return false, nil
if wfStatus.Terminated {
done = true
return
}
if err := wfCtx.Commit(); err != nil {
return false, errors.WithMessage(err, "commit workflow context")
}
wfStatus.StepIndex++
if wfStatus.Terminated || wfStatus.Suspend {
return true, nil
if wfStatus.Suspend {
pause = true
return
}
}
return true, nil // all steps done
return true, false, nil // all steps done
}
+40 -9
View File
@@ -65,9 +65,9 @@ var _ = Describe("Test Workflow", func() {
},
})
wf := NewWorkflow(app, k8sClient)
done, err := wf.ExecuteSteps(context.Background(), revision, runners)
done, pause, err := wf.ExecuteSteps(context.Background(), revision, runners)
Expect(err).ToNot(HaveOccurred())
Expect(pause).Should(BeFalse())
Expect(done).Should(BeFalse())
workflowStatus := app.Status.Workflow
Expect(workflowStatus.ContextBackend.Name).Should(BeEquivalentTo("workflow-" + revision))
@@ -103,8 +103,9 @@ var _ = Describe("Test Workflow", func() {
app.Status.Workflow = workflowStatus
wf = NewWorkflow(app, k8sClient)
done, err = wf.ExecuteSteps(context.Background(), revision, runners)
done, pause, err = wf.ExecuteSteps(context.Background(), revision, runners)
Expect(err).ToNot(HaveOccurred())
Expect(pause).Should(BeFalse())
Expect(done).Should(BeTrue())
app.Status.Workflow.ContextBackend = nil
Expect(cmp.Diff(*app.Status.Workflow, common.WorkflowStatus{
@@ -143,9 +144,10 @@ var _ = Describe("Test Workflow", func() {
},
})
wf := NewWorkflow(app, k8sClient)
done, err := wf.ExecuteSteps(context.Background(), revision, runners)
done, pause, err := wf.ExecuteSteps(context.Background(), revision, runners)
Expect(err).ToNot(HaveOccurred())
Expect(done).Should(BeTrue())
Expect(done).Should(BeFalse())
Expect(pause).Should(BeTrue())
app.Status.Workflow.ContextBackend = nil
Expect(cmp.Diff(*app.Status.Workflow, common.WorkflowStatus{
AppRevision: revision,
@@ -162,9 +164,17 @@ var _ = Describe("Test Workflow", func() {
}},
})).Should(BeEquivalentTo(""))
app.Status.Workflow.Suspend = false
done, err = wf.ExecuteSteps(context.Background(), revision, runners)
// check suspend...
done, pause, err = wf.ExecuteSteps(context.Background(), revision, runners)
Expect(err).ToNot(HaveOccurred())
Expect(pause).Should(BeTrue())
Expect(done).Should(BeFalse())
// check resume
app.Status.Workflow.Suspend = false
done, pause, err = wf.ExecuteSteps(context.Background(), revision, runners)
Expect(err).ToNot(HaveOccurred())
Expect(pause).Should(BeFalse())
Expect(done).Should(BeTrue())
app.Status.Workflow.ContextBackend = nil
Expect(cmp.Diff(*app.Status.Workflow, common.WorkflowStatus{
@@ -184,6 +194,11 @@ var _ = Describe("Test Workflow", func() {
Phase: common.WorkflowStepPhaseSucceeded,
}},
})).Should(BeEquivalentTo(""))
done, pause, err = wf.ExecuteSteps(context.Background(), revision, runners)
Expect(err).ToNot(HaveOccurred())
Expect(pause).Should(BeFalse())
Expect(done).Should(BeTrue())
})
It("test for terminate", func() {
@@ -198,8 +213,9 @@ var _ = Describe("Test Workflow", func() {
},
})
wf := NewWorkflow(app, k8sClient)
done, err := wf.ExecuteSteps(context.Background(), revision, runners)
done, pause, err := wf.ExecuteSteps(context.Background(), revision, runners)
Expect(err).ToNot(HaveOccurred())
Expect(pause).Should(BeFalse())
Expect(done).Should(BeTrue())
app.Status.Workflow.ContextBackend = nil
Expect(cmp.Diff(*app.Status.Workflow, common.WorkflowStatus{
@@ -216,6 +232,11 @@ var _ = Describe("Test Workflow", func() {
Phase: common.WorkflowStepPhaseSucceeded,
}},
})).Should(BeEquivalentTo(""))
done, pause, err = wf.ExecuteSteps(context.Background(), revision, runners)
Expect(err).ToNot(HaveOccurred())
Expect(pause).Should(BeFalse())
Expect(done).Should(BeTrue())
})
It("test for error", func() {
@@ -230,9 +251,10 @@ var _ = Describe("Test Workflow", func() {
},
})
wf := NewWorkflow(app, k8sClient)
done, err := wf.ExecuteSteps(context.Background(), revision, runners)
done, pause, err := wf.ExecuteSteps(context.Background(), revision, runners)
Expect(err).To(HaveOccurred())
Expect(done).Should(BeFalse())
Expect(pause).Should(BeFalse())
app.Status.Workflow.ContextBackend = nil
Expect(cmp.Diff(*app.Status.Workflow, common.WorkflowStatus{
AppRevision: revision,
@@ -244,6 +266,15 @@ var _ = Describe("Test Workflow", func() {
}},
})).Should(BeEquivalentTo(""))
})
It("skip workflow", func() {
app, runners := makeTestCase([]oamcore.WorkflowStep{})
wf := NewWorkflow(app, k8sClient)
done, pause, err := wf.ExecuteSteps(context.Background(), revision, runners)
Expect(err).ToNot(HaveOccurred())
Expect(done).Should(BeTrue())
Expect(pause).Should(BeFalse())
})
})
func makeTestCase(steps []oamcore.WorkflowStep) (*oamcore.Application, []wfTypes.TaskRunner) {