From c463b147a145c2d87f75b3a56fa51b98addad97a Mon Sep 17 00:00:00 2001 From: Hongchao Deng Date: Wed, 2 Jun 2021 22:21:32 -0700 Subject: [PATCH] add PolicyDefinition and WorkflowStepDefinition controller (#1751) --- .../definitionrevision_test.go | 295 ++++++++++++++++++ .../policydefinition_controller.go | 205 ++++++++++++ .../policies/policydefinition/suite_test.go | 121 +++++++ .../definitionrevision_test.go | 295 ++++++++++++++++++ .../workflowstepdefinition/suite_test.go | 121 +++++++ .../workflowstepdefinition_controller.go | 205 ++++++++++++ pkg/controller/core.oam.dev/v1alpha2/setup.go | 4 +- pkg/oam/labels.go | 2 +- 8 files changed, 1246 insertions(+), 2 deletions(-) create mode 100644 pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition/definitionrevision_test.go create mode 100644 pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition/policydefinition_controller.go create mode 100644 pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition/suite_test.go create mode 100644 pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition/definitionrevision_test.go create mode 100644 pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition/suite_test.go create mode 100644 pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition/workflowstepdefinition_controller.go diff --git a/pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition/definitionrevision_test.go b/pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition/definitionrevision_test.go new file mode 100644 index 000000000..e9d02c5d5 --- /dev/null +++ b/pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition/definitionrevision_test.go @@ -0,0 +1,295 @@ +/* + Copyright 2021. The KubeVela Authors. + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. +*/ + +package policydefinition + +import ( + "context" + "fmt" + "time" + + . "github.com/onsi/ginkgo" + . "github.com/onsi/gomega" + v1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/reconcile" + + "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/oam" + "github.com/oam-dev/kubevela/pkg/oam/util" +) + +var _ = Describe("Test DefinitionRevision created by PolicyDefinition", func() { + ctx := context.Background() + namespace := "test-revision" + var ns v1.Namespace + + BeforeEach(func() { + ns = v1.Namespace{ + ObjectMeta: metav1.ObjectMeta{ + Name: namespace, + }, + } + Expect(k8sClient.Create(ctx, &ns)).Should(SatisfyAny(BeNil(), &util.AlreadyExistMatcher{})) + }) + + Context("Test PolicyDefinition", func() { + It("Test update PolicyDefinition", func() { + defName := "test-def" + req := reconcile.Request{NamespacedName: client.ObjectKey{Name: defName, Namespace: namespace}} + + def1 := defWithNoTemplate.DeepCopy() + def1.Name = defName + def1.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, "test-v1") + By("create policyDefinition") + Expect(k8sClient.Create(ctx, def1)).Should(SatisfyAll(BeNil())) + reconcileRetry(&r, req) + + By("check whether definitionRevision is created") + defRevName1 := fmt.Sprintf("%s-v1", defName) + var defRev1 v1beta1.DefinitionRevision + Eventually(func() error { + err := k8sClient.Get(ctx, client.ObjectKey{Namespace: namespace, Name: defRevName1}, &defRev1) + return err + }, 10*time.Second, time.Second).Should(BeNil()) + + By("update policyDefinition") + def := new(v1beta1.PolicyDefinition) + Eventually(func() error { + err := k8sClient.Get(ctx, client.ObjectKey{Namespace: namespace, Name: defName}, def) + if err != nil { + return err + } + def.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, "test-v2") + return k8sClient.Update(ctx, def) + }, 10*time.Second, time.Second).Should(BeNil()) + + reconcileRetry(&r, req) + + By("check whether a new definitionRevision is created") + tdRevName2 := fmt.Sprintf("%s-v2", defName) + var tdRev2 v1beta1.DefinitionRevision + Eventually(func() error { + return k8sClient.Get(ctx, client.ObjectKey{Namespace: namespace, Name: tdRevName2}, &tdRev2) + }, 10*time.Second, time.Second).Should(BeNil()) + }) + + It("Test only update PolicyDefinition Labels, Shouldn't create new revision", func() { + def := defWithNoTemplate.DeepCopy() + defName := "test-def-update" + def.Name = defName + defKey := client.ObjectKey{Namespace: namespace, Name: defName} + req := reconcile.Request{NamespacedName: defKey} + def.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, "test") + Expect(k8sClient.Create(ctx, def)).Should(BeNil()) + reconcileRetry(&r, req) + + By("Check revision create by PolicyDefinition") + defRevName := fmt.Sprintf("%s-v1", defName) + revKey := client.ObjectKey{Namespace: namespace, Name: defRevName} + var defRev v1beta1.DefinitionRevision + Eventually(func() error { + return k8sClient.Get(ctx, revKey, &defRev) + }, 10*time.Second, time.Second).Should(BeNil()) + + By("Only update PolicyDefinition Labels") + var checkRev v1beta1.PolicyDefinition + Eventually(func() error { + err := k8sClient.Get(ctx, defKey, &checkRev) + if err != nil { + return err + } + checkRev.SetLabels(map[string]string{ + "test-label": "test-defRev", + }) + return k8sClient.Update(ctx, &checkRev) + }, 10*time.Second, time.Second).Should(BeNil()) + reconcileRetry(&r, req) + + newDefRevName := fmt.Sprintf("%s-v2", defName) + newRevKey := client.ObjectKey{Namespace: namespace, Name: newDefRevName} + Expect(k8sClient.Get(ctx, newRevKey, &defRev)).Should(HaveOccurred()) + }) + }) + + Context("Test PolicyDefinition Controller clean up", func() { + It("Test clean up definitionRevision", func() { + var revKey client.ObjectKey + var defRev v1beta1.DefinitionRevision + defName := "test-clean-up" + revisionNum := 1 + defKey := client.ObjectKey{Namespace: namespace, Name: defName} + req := reconcile.Request{NamespacedName: defKey} + + By("create a new PolicyDefinition") + def := defWithNoTemplate.DeepCopy() + def.Name = defName + def.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, fmt.Sprintf("test-v%d", revisionNum)) + Expect(k8sClient.Create(ctx, def)).Should(BeNil()) + + By("update PolicyDefinition") + checkRev := new(v1beta1.PolicyDefinition) + for i := 0; i < defRevisionLimit+1; i++ { + Eventually(func() error { + err := k8sClient.Get(ctx, defKey, checkRev) + if err != nil { + return err + } + checkRev.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, fmt.Sprintf("test-v%d", revisionNum)) + return k8sClient.Update(ctx, checkRev) + }, 10*time.Second, time.Second).Should(BeNil()) + reconcileRetry(&r, req) + + revKey = client.ObjectKey{Namespace: namespace, Name: fmt.Sprintf("%s-v%d", defName, revisionNum)} + revisionNum++ + var defRev v1beta1.DefinitionRevision + Eventually(func() error { + return k8sClient.Get(ctx, revKey, &defRev) + }, 10*time.Second, time.Second).Should(BeNil()) + } + + By("create new PolicyDefinition will remove oldest definitionRevision") + Eventually(func() error { + err := k8sClient.Get(ctx, defKey, checkRev) + if err != nil { + return err + } + checkRev.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, fmt.Sprintf("test-v%d", revisionNum)) + return k8sClient.Update(ctx, checkRev) + }, 10*time.Second, time.Second).Should(BeNil()) + reconcileRetry(&r, req) + + revKey = client.ObjectKey{Namespace: namespace, Name: fmt.Sprintf("%s-v%d", defName, revisionNum)} + revisionNum++ + Eventually(func() error { + return k8sClient.Get(ctx, revKey, &defRev) + }, 10*time.Second, time.Second).Should(BeNil()) + + deletedRevision := new(v1beta1.DefinitionRevision) + deleteRevKey := types.NamespacedName{Namespace: namespace, Name: defName + "-v1"} + listOpts := []client.ListOption{ + client.InNamespace(namespace), + client.MatchingLabels{ + oam.LabelPolicyDefinitionName: defName, + }, + } + defRevList := new(v1beta1.DefinitionRevisionList) + Eventually(func() error { + err := k8sClient.List(ctx, defRevList, listOpts...) + if err != nil { + return err + } + if len(defRevList.Items) != defRevisionLimit+1 { + return fmt.Errorf("error defRevison number wants %d, actually %d", defRevisionLimit+1, len(defRevList.Items)) + } + err = k8sClient.Get(ctx, deleteRevKey, deletedRevision) + if err == nil || !apierrors.IsNotFound(err) { + return fmt.Errorf("haven't clean up the oldest revision") + } + return nil + }, time.Second*30, time.Microsecond*300).Should(BeNil()) + + By("update app again will continue to delete the oldest revision") + Eventually(func() error { + err := k8sClient.Get(ctx, defKey, checkRev) + if err != nil { + return err + } + checkRev.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, fmt.Sprintf("test-v%d", revisionNum)) + return k8sClient.Update(ctx, checkRev) + }, 10*time.Second, time.Second).Should(BeNil()) + reconcileRetry(&r, req) + + revKey = client.ObjectKey{Namespace: namespace, Name: fmt.Sprintf("%s-v%d", defName, revisionNum)} + Eventually(func() error { + return k8sClient.Get(ctx, revKey, &defRev) + }, 10*time.Second, time.Second).Should(BeNil()) + + deleteRevKey = types.NamespacedName{Namespace: namespace, Name: defName + "-v2"} + Eventually(func() error { + err := k8sClient.List(ctx, defRevList, listOpts...) + if err != nil { + return err + } + if len(defRevList.Items) != defRevisionLimit+1 { + return fmt.Errorf("error defRevison number wants %d, actually %d", defRevisionLimit+1, len(defRevList.Items)) + } + err = k8sClient.Get(ctx, deleteRevKey, deletedRevision) + if err == nil || !apierrors.IsNotFound(err) { + return fmt.Errorf("haven't clean up the oldest revision") + } + return nil + }, time.Second*30, time.Microsecond*300).Should(BeNil()) + }) + }) +}) + +var defTemplate = ` +output: { + apiVersion: "batch/v1" + kind: "Job" + spec: { + parallelism: parameter.count + completions: parameter.count + template: spec: { + restartPolicy: parameter.restart + containers: [{ + name: "%s" + image: parameter.image + + if parameter["cmd"] != _|_ { + command: parameter.cmd + } + }] + } + } +} +parameter: { + // +usage=Specify number of tasks to run in parallel + // +short=c + count: *1 | int + + // +usage=Which image would you like to use for your service + // +short=i + image: string + + // +usage=Define the job restart policy, the value can only be Never or OnFailure. By default, it's Never. + restart: *"Never" | string + + // +usage=Commands to run in the container + cmd?: [...string] +} +` + +var defWithNoTemplate = &v1beta1.PolicyDefinition{ + TypeMeta: metav1.TypeMeta{ + Kind: "PolicyDefinition", + APIVersion: "core.oam.dev/v1beta1", + }, + ObjectMeta: metav1.ObjectMeta{ + Name: "test-defrev", + Namespace: "test-revision", + }, + Spec: v1beta1.PolicyDefinitionSpec{ + Schematic: &common.Schematic{ + CUE: &common.CUE{}, + }, + }, +} diff --git a/pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition/policydefinition_controller.go b/pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition/policydefinition_controller.go new file mode 100644 index 000000000..fd0aed497 --- /dev/null +++ b/pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition/policydefinition_controller.go @@ -0,0 +1,205 @@ +/* + + Copyright 2021 The KubeVela Authors. + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. + +*/ + +package policydefinition + +import ( + "context" + "fmt" + + cpv1alpha1 "github.com/crossplane/crossplane-runtime/apis/core/v1alpha1" + "github.com/crossplane/crossplane-runtime/pkg/event" + "github.com/crossplane/crossplane-runtime/pkg/logging" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/client-go/util/retry" + "k8s.io/klog/v2" + "k8s.io/utils/pointer" + ctrl "sigs.k8s.io/controller-runtime" + "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" + controller "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev" + coredef "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/core" + "github.com/oam-dev/kubevela/pkg/controller/utils" + "github.com/oam-dev/kubevela/pkg/cue/packages" + "github.com/oam-dev/kubevela/pkg/oam" + "github.com/oam-dev/kubevela/pkg/oam/discoverymapper" + "github.com/oam-dev/kubevela/pkg/oam/util" +) + +// Reconciler reconciles a PolicyDefinition object +type Reconciler struct { + client.Client + dm discoverymapper.DiscoveryMapper + pd *packages.PackageDiscover + Scheme *runtime.Scheme + record event.Recorder + defRevLimit int +} + +// Reconcile is the main logic for PolicyDefinition controller +func (r *Reconciler) Reconcile(req ctrl.Request) (ctrl.Result, error) { + definitionName := req.NamespacedName.Name + klog.InfoS("Reconciling PolicyDefinition...", "Name", definitionName, "Namespace", req.Namespace) + ctx := context.Background() + + var policydefinition v1beta1.PolicyDefinition + if err := r.Get(ctx, req.NamespacedName, &policydefinition); err != nil { + if apierrors.IsNotFound(err) { + err = nil + } + return ctrl.Result{}, err + } + + // this is a placeholder for finalizer here in the future + if policydefinition.DeletionTimestamp != nil { + return ctrl.Result{}, nil + } + + // refresh package discover when policyDefinition is registered + if policydefinition.Spec.Reference.Name != "" { + err := utils.RefreshPackageDiscover(ctx, r.Client, r.dm, r.pd, &policydefinition) + if err != nil { + klog.ErrorS(err, "cannot refresh packageDiscover") + r.record.Event(&policydefinition, event.Warning("cannot refresh packageDiscover", err)) + return ctrl.Result{}, util.PatchCondition(ctx, r, &policydefinition, + cpv1alpha1.ReconcileError(fmt.Errorf(util.ErrRefreshPackageDiscover, err))) + } + } + + // generate DefinitionRevision from policyDefinition + defRev, isNewRevision, err := coredef.GenerateDefinitionRevision(ctx, r.Client, &policydefinition) + if err != nil { + klog.ErrorS(err, "cannot generate DefinitionRevision", "PolicyDefinitionName", policydefinition.Name) + r.record.Event(&policydefinition, event.Warning("cannot generate DefinitionRevision", err)) + return ctrl.Result{}, util.PatchCondition(ctx, r, &policydefinition, + cpv1alpha1.ReconcileError(fmt.Errorf(util.ErrGenerateDefinitionRevision, policydefinition.Name, err))) + } + if !isNewRevision { + if err = r.createOrUpdatePolicyDefRevision(ctx, req.Namespace, &policydefinition, defRev); err != nil { + klog.ErrorS(err, "cannot update DefinitionRevision") + r.record.Event(&(policydefinition), event.Warning("cannot update DefinitionRevision", err)) + return ctrl.Result{}, util.PatchCondition(ctx, r, &(policydefinition), + cpv1alpha1.ReconcileError(fmt.Errorf(util.ErrCreateOrUpdateDefinitionRevision, defRev.Name, err))) + } + klog.InfoS("Successfully update DefinitionRevision", "name", defRev.Name) + + if err := coredef.CleanUpDefinitionRevision(ctx, r.Client, &policydefinition, r.defRevLimit); err != nil { + klog.Error("[Garbage collection]") + r.record.Event(&policydefinition, event.Warning("failed to garbage collect DefinitionRevision of type PolicyDefinition", err)) + } + return ctrl.Result{}, nil + } + + if err = r.createOrUpdatePolicyDefRevision(ctx, req.Namespace, &policydefinition, defRev); err != nil { + klog.ErrorS(err, "cannot create DefinitionRevision") + r.record.Event(&(policydefinition), event.Warning("cannot create DefinitionRevision", err)) + return ctrl.Result{}, util.PatchCondition(ctx, r, &(policydefinition), + cpv1alpha1.ReconcileError(fmt.Errorf(util.ErrCreateOrUpdateDefinitionRevision, defRev.Name, err))) + } + klog.InfoS("Successfully createOrUpdatePolicyDefRevision", "name", defRev.Name) + + policydefinition.Status.LatestRevision = &common.Revision{ + Name: defRev.Name, + Revision: defRev.Spec.Revision, + RevisionHash: defRev.Spec.RevisionHash, + } + + if err := r.UpdateStatus(ctx, &policydefinition); err != nil { + klog.ErrorS(err, "cannot update PolicyDefinition Status") + r.record.Event(&(policydefinition), event.Warning("cannot update PolicyDefinition Status", err)) + return ctrl.Result{}, util.PatchCondition(ctx, r, &(policydefinition), + cpv1alpha1.ReconcileError(fmt.Errorf(util.ErrUpdatePolicyDefinition, policydefinition.Name, err))) + } + + if err := coredef.CleanUpDefinitionRevision(ctx, r.Client, &policydefinition, r.defRevLimit); err != nil { + klog.Error("[Garbage collection]") + r.record.Event(&policydefinition, event.Warning("failed to garbage collect DefinitionRevision of type PolicyDefinition", err)) + } + + return ctrl.Result{}, nil +} + +func (r *Reconciler) createOrUpdatePolicyDefRevision(ctx context.Context, ns string, + def *v1beta1.PolicyDefinition, defRev *v1beta1.DefinitionRevision) error { + + ownerReference := []metav1.OwnerReference{{ + APIVersion: def.APIVersion, + Kind: def.Kind, + Name: def.Name, + UID: def.GetUID(), + Controller: pointer.BoolPtr(true), + BlockOwnerDeletion: pointer.BoolPtr(true), + }} + + defRev.SetLabels(def.GetLabels()) + defRev.SetLabels(util.MergeMapOverrideWithDst(defRev.Labels, + map[string]string{oam.LabelPolicyDefinitionName: def.Name})) + defRev.SetNamespace(ns) + defRev.SetAnnotations(def.GetAnnotations()) + defRev.SetOwnerReferences(ownerReference) + + rev := &v1beta1.DefinitionRevision{} + if err := r.Client.Get(ctx, client.ObjectKey{Namespace: ns, Name: defRev.Name}, rev); err != nil { + if apierrors.IsNotFound(err) { + return r.Create(ctx, defRev) + } + return err + } + + rev.SetAnnotations(defRev.GetAnnotations()) + rev.SetLabels(defRev.GetLabels()) + rev.SetOwnerReferences(ownerReference) + return r.Update(ctx, rev) +} + +// UpdateStatus updates v1beta1.PolicyDefinition's Status with retry.RetryOnConflict +func (r *Reconciler) UpdateStatus(ctx context.Context, def *v1beta1.PolicyDefinition, opts ...client.UpdateOption) error { + status := def.DeepCopy().Status + return retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) { + if err = r.Get(ctx, client.ObjectKey{Namespace: def.Namespace, Name: def.Name}, def); err != nil { + return + } + def.Status = status + return r.Status().Update(ctx, def, opts...) + }) +} + +// SetupWithManager will setup with event recorder +func (r *Reconciler) SetupWithManager(mgr ctrl.Manager) error { + r.record = event.NewAPIRecorder(mgr.GetEventRecorderFor("PolicyDefinition")). + WithAnnotations("controller", "PolicyDefinition") + return ctrl.NewControllerManagedBy(mgr). + For(&v1beta1.PolicyDefinition{}). + Complete(r) +} + +// Setup adds a controller that reconciles PolicyDefinition. +func Setup(mgr ctrl.Manager, args controller.Args, _ logging.Logger) error { + r := Reconciler{ + Client: mgr.GetClient(), + Scheme: mgr.GetScheme(), + dm: args.DiscoveryMapper, + pd: args.PackageDiscover, + defRevLimit: args.DefRevisionLimit, + } + return r.SetupWithManager(mgr) +} diff --git a/pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition/suite_test.go b/pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition/suite_test.go new file mode 100644 index 000000000..612eb6c34 --- /dev/null +++ b/pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition/suite_test.go @@ -0,0 +1,121 @@ +/* + Copyright 2021 The KubeVela Authors. + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. + +*/ + +package policydefinition + +import ( + "path/filepath" + "testing" + "time" + + . "github.com/onsi/ginkgo" + . "github.com/onsi/gomega" + + crdv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" + "k8s.io/client-go/kubernetes/scheme" + "k8s.io/client-go/rest" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/envtest" + "sigs.k8s.io/controller-runtime/pkg/reconcile" + + oamCore "github.com/oam-dev/kubevela/apis/core.oam.dev" + "github.com/oam-dev/kubevela/pkg/cue/packages" + "github.com/oam-dev/kubevela/pkg/oam/discoverymapper" +) + +var cfg *rest.Config +var k8sClient client.Client +var testEnv *envtest.Environment +var controllerDone chan struct{} +var r Reconciler +var defRevisionLimit = 5 + +func TestPolicyDefinition(t *testing.T) { + RegisterFailHandler(Fail) + RunSpecs(t, "PolicyDefinition Suite") +} + +var _ = BeforeSuite(func(done Done) { + By("Bootstrapping test environment") + useExistCluster := false + testEnv = &envtest.Environment{ + CRDDirectoryPaths: []string{ + filepath.Join("../../../../../../..", "charts/vela-core/crds"), // this has all the required CRDs, + }, + UseExistingCluster: &useExistCluster, + } + var err error + cfg, err = testEnv.Start() + Expect(err).ToNot(HaveOccurred()) + Expect(cfg).ToNot(BeNil()) + + Expect(oamCore.AddToScheme(scheme.Scheme)).Should(BeNil()) + Expect(crdv1.AddToScheme(scheme.Scheme)).Should(BeNil()) + + By("Create the k8s client") + k8sClient, err = client.New(cfg, client.Options{Scheme: scheme.Scheme}) + Expect(err).ToNot(HaveOccurred()) + Expect(k8sClient).ToNot(BeNil()) + + By("Starting the controller in the background") + mgr, err := ctrl.NewManager(cfg, ctrl.Options{ + Scheme: scheme.Scheme, + MetricsBindAddress: "0", + Port: 48081, + }) + Expect(err).ToNot(HaveOccurred()) + + pd, err := packages.NewPackageDiscover(cfg) + Expect(err).ToNot(HaveOccurred()) + dm, err := discoverymapper.New(mgr.GetConfig()) + Expect(err).ToNot(HaveOccurred()) + _, err = dm.Refresh() + Expect(err).ToNot(HaveOccurred()) + + r = Reconciler{ + Client: mgr.GetClient(), + Scheme: mgr.GetScheme(), + dm: dm, + pd: pd, + defRevLimit: defRevisionLimit, + } + Expect(r.SetupWithManager(mgr)).ToNot(HaveOccurred()) + controllerDone = make(chan struct{}, 1) + go func() { + defer GinkgoRecover() + Expect(mgr.Start(controllerDone)).ToNot(HaveOccurred()) + }() + + close(done) +}, 60) + +var _ = AfterSuite(func() { + By("Stop the controller") + close(controllerDone) + + By("Tearing down the test environment") + err := testEnv.Stop() + Expect(err).ToNot(HaveOccurred()) +}) + +func reconcileRetry(r reconcile.Reconciler, req reconcile.Request) { + Eventually(func() error { + _, err := r.Reconcile(req) + return err + }, 15*time.Second, time.Second).Should(BeNil()) +} diff --git a/pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition/definitionrevision_test.go b/pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition/definitionrevision_test.go new file mode 100644 index 000000000..4d8e06614 --- /dev/null +++ b/pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition/definitionrevision_test.go @@ -0,0 +1,295 @@ +/* + Copyright 2021. The KubeVela Authors. + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. +*/ + +package workflowstepdefinition + +import ( + "context" + "fmt" + "time" + + . "github.com/onsi/ginkgo" + . "github.com/onsi/gomega" + v1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/reconcile" + + "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/oam" + "github.com/oam-dev/kubevela/pkg/oam/util" +) + +var _ = Describe("Test DefinitionRevision created by WorkflowStepDefinition", func() { + ctx := context.Background() + namespace := "test-revision" + var ns v1.Namespace + + BeforeEach(func() { + ns = v1.Namespace{ + ObjectMeta: metav1.ObjectMeta{ + Name: namespace, + }, + } + Expect(k8sClient.Create(ctx, &ns)).Should(SatisfyAny(BeNil(), &util.AlreadyExistMatcher{})) + }) + + Context("Test WorkflowStepDefinition", func() { + It("Test update WorkflowStepDefinition", func() { + defName := "test-def" + req := reconcile.Request{NamespacedName: client.ObjectKey{Name: defName, Namespace: namespace}} + + def1 := defWithNoTemplate.DeepCopy() + def1.Name = defName + def1.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, "test-v1") + By("create workflowStepDefinition") + Expect(k8sClient.Create(ctx, def1)).Should(SatisfyAll(BeNil())) + reconcileRetry(&r, req) + + By("check whether definitionRevision is created") + defRevName1 := fmt.Sprintf("%s-v1", defName) + var defRev1 v1beta1.DefinitionRevision + Eventually(func() error { + err := k8sClient.Get(ctx, client.ObjectKey{Namespace: namespace, Name: defRevName1}, &defRev1) + return err + }, 10*time.Second, time.Second).Should(BeNil()) + + By("update workflowStepDefinition") + def := new(v1beta1.WorkflowStepDefinition) + Eventually(func() error { + err := k8sClient.Get(ctx, client.ObjectKey{Namespace: namespace, Name: defName}, def) + if err != nil { + return err + } + def.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, "test-v2") + return k8sClient.Update(ctx, def) + }, 10*time.Second, time.Second).Should(BeNil()) + + reconcileRetry(&r, req) + + By("check whether a new definitionRevision is created") + tdRevName2 := fmt.Sprintf("%s-v2", defName) + var tdRev2 v1beta1.DefinitionRevision + Eventually(func() error { + return k8sClient.Get(ctx, client.ObjectKey{Namespace: namespace, Name: tdRevName2}, &tdRev2) + }, 10*time.Second, time.Second).Should(BeNil()) + }) + + It("Test only update WorkflowStepDefinition Labels, Shouldn't create new revision", func() { + def := defWithNoTemplate.DeepCopy() + defName := "test-def-update" + def.Name = defName + defKey := client.ObjectKey{Namespace: namespace, Name: defName} + req := reconcile.Request{NamespacedName: defKey} + def.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, "test") + Expect(k8sClient.Create(ctx, def)).Should(BeNil()) + reconcileRetry(&r, req) + + By("Check revision create by WorkflowStepDefinition") + defRevName := fmt.Sprintf("%s-v1", defName) + revKey := client.ObjectKey{Namespace: namespace, Name: defRevName} + var defRev v1beta1.DefinitionRevision + Eventually(func() error { + return k8sClient.Get(ctx, revKey, &defRev) + }, 10*time.Second, time.Second).Should(BeNil()) + + By("Only update WorkflowStepDefinition Labels") + var checkRev v1beta1.WorkflowStepDefinition + Eventually(func() error { + err := k8sClient.Get(ctx, defKey, &checkRev) + if err != nil { + return err + } + checkRev.SetLabels(map[string]string{ + "test-label": "test-defRev", + }) + return k8sClient.Update(ctx, &checkRev) + }, 10*time.Second, time.Second).Should(BeNil()) + reconcileRetry(&r, req) + + newDefRevName := fmt.Sprintf("%s-v2", defName) + newRevKey := client.ObjectKey{Namespace: namespace, Name: newDefRevName} + Expect(k8sClient.Get(ctx, newRevKey, &defRev)).Should(HaveOccurred()) + }) + }) + + Context("Test WorkflowStepDefinition Controller clean up", func() { + It("Test clean up definitionRevision", func() { + var revKey client.ObjectKey + var defRev v1beta1.DefinitionRevision + defName := "test-clean-up" + revisionNum := 1 + defKey := client.ObjectKey{Namespace: namespace, Name: defName} + req := reconcile.Request{NamespacedName: defKey} + + By("create a new WorkflowStepDefinition") + def := defWithNoTemplate.DeepCopy() + def.Name = defName + def.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, fmt.Sprintf("test-v%d", revisionNum)) + Expect(k8sClient.Create(ctx, def)).Should(BeNil()) + + By("update WorkflowStepDefinition") + checkRev := new(v1beta1.WorkflowStepDefinition) + for i := 0; i < defRevisionLimit+1; i++ { + Eventually(func() error { + err := k8sClient.Get(ctx, defKey, checkRev) + if err != nil { + return err + } + checkRev.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, fmt.Sprintf("test-v%d", revisionNum)) + return k8sClient.Update(ctx, checkRev) + }, 10*time.Second, time.Second).Should(BeNil()) + reconcileRetry(&r, req) + + revKey = client.ObjectKey{Namespace: namespace, Name: fmt.Sprintf("%s-v%d", defName, revisionNum)} + revisionNum++ + var defRev v1beta1.DefinitionRevision + Eventually(func() error { + return k8sClient.Get(ctx, revKey, &defRev) + }, 10*time.Second, time.Second).Should(BeNil()) + } + + By("create new WorkflowStepDefinition will remove oldest definitionRevision") + Eventually(func() error { + err := k8sClient.Get(ctx, defKey, checkRev) + if err != nil { + return err + } + checkRev.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, fmt.Sprintf("test-v%d", revisionNum)) + return k8sClient.Update(ctx, checkRev) + }, 10*time.Second, time.Second).Should(BeNil()) + reconcileRetry(&r, req) + + revKey = client.ObjectKey{Namespace: namespace, Name: fmt.Sprintf("%s-v%d", defName, revisionNum)} + revisionNum++ + Eventually(func() error { + return k8sClient.Get(ctx, revKey, &defRev) + }, 10*time.Second, time.Second).Should(BeNil()) + + deletedRevision := new(v1beta1.DefinitionRevision) + deleteRevKey := types.NamespacedName{Namespace: namespace, Name: defName + "-v1"} + listOpts := []client.ListOption{ + client.InNamespace(namespace), + client.MatchingLabels{ + oam.LabelWorkflowStepDefinitionName: defName, + }, + } + defRevList := new(v1beta1.DefinitionRevisionList) + Eventually(func() error { + err := k8sClient.List(ctx, defRevList, listOpts...) + if err != nil { + return err + } + if len(defRevList.Items) != defRevisionLimit+1 { + return fmt.Errorf("error defRevison number wants %d, actually %d", defRevisionLimit+1, len(defRevList.Items)) + } + err = k8sClient.Get(ctx, deleteRevKey, deletedRevision) + if err == nil || !apierrors.IsNotFound(err) { + return fmt.Errorf("haven't clean up the oldest revision") + } + return nil + }, time.Second*30, time.Microsecond*300).Should(BeNil()) + + By("update app again will continue to delete the oldest revision") + Eventually(func() error { + err := k8sClient.Get(ctx, defKey, checkRev) + if err != nil { + return err + } + checkRev.Spec.Schematic.CUE.Template = fmt.Sprintf(defTemplate, fmt.Sprintf("test-v%d", revisionNum)) + return k8sClient.Update(ctx, checkRev) + }, 10*time.Second, time.Second).Should(BeNil()) + reconcileRetry(&r, req) + + revKey = client.ObjectKey{Namespace: namespace, Name: fmt.Sprintf("%s-v%d", defName, revisionNum)} + Eventually(func() error { + return k8sClient.Get(ctx, revKey, &defRev) + }, 10*time.Second, time.Second).Should(BeNil()) + + deleteRevKey = types.NamespacedName{Namespace: namespace, Name: defName + "-v2"} + Eventually(func() error { + err := k8sClient.List(ctx, defRevList, listOpts...) + if err != nil { + return err + } + if len(defRevList.Items) != defRevisionLimit+1 { + return fmt.Errorf("error defRevison number wants %d, actually %d", defRevisionLimit+1, len(defRevList.Items)) + } + err = k8sClient.Get(ctx, deleteRevKey, deletedRevision) + if err == nil || !apierrors.IsNotFound(err) { + return fmt.Errorf("haven't clean up the oldest revision") + } + return nil + }, time.Second*30, time.Microsecond*300).Should(BeNil()) + }) + }) +}) + +var defTemplate = ` +output: { + apiVersion: "batch/v1" + kind: "Job" + spec: { + parallelism: parameter.count + completions: parameter.count + template: spec: { + restartPolicy: parameter.restart + containers: [{ + name: "%s" + image: parameter.image + + if parameter["cmd"] != _|_ { + command: parameter.cmd + } + }] + } + } +} +parameter: { + // +usage=Specify number of tasks to run in parallel + // +short=c + count: *1 | int + + // +usage=Which image would you like to use for your service + // +short=i + image: string + + // +usage=Define the job restart policy, the value can only be Never or OnFailure. By default, it's Never. + restart: *"Never" | string + + // +usage=Commands to run in the container + cmd?: [...string] +} +` + +var defWithNoTemplate = &v1beta1.WorkflowStepDefinition{ + TypeMeta: metav1.TypeMeta{ + Kind: "WorkflowStepDefinition", + APIVersion: "core.oam.dev/v1beta1", + }, + ObjectMeta: metav1.ObjectMeta{ + Name: "test-defrev", + Namespace: "test-revision", + }, + Spec: v1beta1.WorkflowStepDefinitionSpec{ + Schematic: &common.Schematic{ + CUE: &common.CUE{}, + }, + }, +} diff --git a/pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition/suite_test.go b/pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition/suite_test.go new file mode 100644 index 000000000..73d34fff2 --- /dev/null +++ b/pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition/suite_test.go @@ -0,0 +1,121 @@ +/* + Copyright 2021 The KubeVela Authors. + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. + +*/ + +package workflowstepdefinition + +import ( + "path/filepath" + "testing" + "time" + + . "github.com/onsi/ginkgo" + . "github.com/onsi/gomega" + + crdv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" + "k8s.io/client-go/kubernetes/scheme" + "k8s.io/client-go/rest" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/envtest" + "sigs.k8s.io/controller-runtime/pkg/reconcile" + + oamCore "github.com/oam-dev/kubevela/apis/core.oam.dev" + "github.com/oam-dev/kubevela/pkg/cue/packages" + "github.com/oam-dev/kubevela/pkg/oam/discoverymapper" +) + +var cfg *rest.Config +var k8sClient client.Client +var testEnv *envtest.Environment +var controllerDone chan struct{} +var r Reconciler +var defRevisionLimit = 5 + +func TestWorkflowStepDefinition(t *testing.T) { + RegisterFailHandler(Fail) + RunSpecs(t, "WorkflowStepDefinition Suite") +} + +var _ = BeforeSuite(func(done Done) { + By("Bootstrapping test environment") + useExistCluster := false + testEnv = &envtest.Environment{ + CRDDirectoryPaths: []string{ + filepath.Join("../../../../../../..", "charts/vela-core/crds"), // this has all the required CRDs, + }, + UseExistingCluster: &useExistCluster, + } + var err error + cfg, err = testEnv.Start() + Expect(err).ToNot(HaveOccurred()) + Expect(cfg).ToNot(BeNil()) + + Expect(oamCore.AddToScheme(scheme.Scheme)).Should(BeNil()) + Expect(crdv1.AddToScheme(scheme.Scheme)).Should(BeNil()) + + By("Create the k8s client") + k8sClient, err = client.New(cfg, client.Options{Scheme: scheme.Scheme}) + Expect(err).ToNot(HaveOccurred()) + Expect(k8sClient).ToNot(BeNil()) + + By("Starting the controller in the background") + mgr, err := ctrl.NewManager(cfg, ctrl.Options{ + Scheme: scheme.Scheme, + MetricsBindAddress: "0", + Port: 48081, + }) + Expect(err).ToNot(HaveOccurred()) + + pd, err := packages.NewPackageDiscover(cfg) + Expect(err).ToNot(HaveOccurred()) + dm, err := discoverymapper.New(mgr.GetConfig()) + Expect(err).ToNot(HaveOccurred()) + _, err = dm.Refresh() + Expect(err).ToNot(HaveOccurred()) + + r = Reconciler{ + Client: mgr.GetClient(), + Scheme: mgr.GetScheme(), + dm: dm, + pd: pd, + defRevLimit: defRevisionLimit, + } + Expect(r.SetupWithManager(mgr)).ToNot(HaveOccurred()) + controllerDone = make(chan struct{}, 1) + go func() { + defer GinkgoRecover() + Expect(mgr.Start(controllerDone)).ToNot(HaveOccurred()) + }() + + close(done) +}, 60) + +var _ = AfterSuite(func() { + By("Stop the controller") + close(controllerDone) + + By("Tearing down the test environment") + err := testEnv.Stop() + Expect(err).ToNot(HaveOccurred()) +}) + +func reconcileRetry(r reconcile.Reconciler, req reconcile.Request) { + Eventually(func() error { + _, err := r.Reconcile(req) + return err + }, 15*time.Second, time.Second).Should(BeNil()) +} diff --git a/pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition/workflowstepdefinition_controller.go b/pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition/workflowstepdefinition_controller.go new file mode 100644 index 000000000..984f5713a --- /dev/null +++ b/pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition/workflowstepdefinition_controller.go @@ -0,0 +1,205 @@ +/* + + Copyright 2021 The KubeVela Authors. + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. + +*/ + +package workflowstepdefinition + +import ( + "context" + "fmt" + + cpv1alpha1 "github.com/crossplane/crossplane-runtime/apis/core/v1alpha1" + "github.com/crossplane/crossplane-runtime/pkg/event" + "github.com/crossplane/crossplane-runtime/pkg/logging" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/client-go/util/retry" + "k8s.io/klog/v2" + "k8s.io/utils/pointer" + ctrl "sigs.k8s.io/controller-runtime" + "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" + controller "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev" + coredef "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/core" + "github.com/oam-dev/kubevela/pkg/controller/utils" + "github.com/oam-dev/kubevela/pkg/cue/packages" + "github.com/oam-dev/kubevela/pkg/oam" + "github.com/oam-dev/kubevela/pkg/oam/discoverymapper" + "github.com/oam-dev/kubevela/pkg/oam/util" +) + +// Reconciler reconciles a WorkflowStepDefinition object +type Reconciler struct { + client.Client + dm discoverymapper.DiscoveryMapper + pd *packages.PackageDiscover + Scheme *runtime.Scheme + record event.Recorder + defRevLimit int +} + +// Reconcile is the main logic for WorkflowStepDefinition controller +func (r *Reconciler) Reconcile(req ctrl.Request) (ctrl.Result, error) { + definitionName := req.NamespacedName.Name + klog.InfoS("Reconciling WorkflowStepDefinition...", "Name", definitionName, "Namespace", req.Namespace) + ctx := context.Background() + + var wfstepdefinition v1beta1.WorkflowStepDefinition + if err := r.Get(ctx, req.NamespacedName, &wfstepdefinition); err != nil { + if apierrors.IsNotFound(err) { + err = nil + } + return ctrl.Result{}, err + } + + // this is a placeholder for finalizer here in the future + if wfstepdefinition.DeletionTimestamp != nil { + return ctrl.Result{}, nil + } + + // refresh package discover when WorkflowStepDefinition is registered + if wfstepdefinition.Spec.Reference.Name != "" { + err := utils.RefreshPackageDiscover(ctx, r.Client, r.dm, r.pd, &wfstepdefinition) + if err != nil { + klog.ErrorS(err, "cannot refresh packageDiscover") + r.record.Event(&wfstepdefinition, event.Warning("cannot refresh packageDiscover", err)) + return ctrl.Result{}, util.PatchCondition(ctx, r, &wfstepdefinition, + cpv1alpha1.ReconcileError(fmt.Errorf(util.ErrRefreshPackageDiscover, err))) + } + } + // generate DefinitionRevision from WorkflowStepDefinition + defRev, isNewRevision, err := coredef.GenerateDefinitionRevision(ctx, r.Client, &wfstepdefinition) + if err != nil { + klog.ErrorS(err, "cannot generate DefinitionRevision", "WorkflowStepDefinitionName", wfstepdefinition.Name) + r.record.Event(&wfstepdefinition, event.Warning("cannot generate DefinitionRevision", err)) + return ctrl.Result{}, util.PatchCondition(ctx, r, &wfstepdefinition, + cpv1alpha1.ReconcileError(fmt.Errorf(util.ErrGenerateDefinitionRevision, wfstepdefinition.Name, err))) + } + + if !isNewRevision { + if err = r.createOrUpdateWFStepDefRevision(ctx, req.Namespace, &wfstepdefinition, defRev); err != nil { + klog.ErrorS(err, "cannot update DefinitionRevision") + r.record.Event(&(wfstepdefinition), event.Warning("cannot update DefinitionRevision", err)) + return ctrl.Result{}, util.PatchCondition(ctx, r, &(wfstepdefinition), + cpv1alpha1.ReconcileError(fmt.Errorf(util.ErrCreateOrUpdateDefinitionRevision, defRev.Name, err))) + } + klog.InfoS("Successfully update DefinitionRevision", "name", defRev.Name) + + if err := coredef.CleanUpDefinitionRevision(ctx, r.Client, &wfstepdefinition, r.defRevLimit); err != nil { + klog.Error("[Garbage collection]") + r.record.Event(&wfstepdefinition, event.Warning("failed to garbage collect DefinitionRevision of type WorkflowStepDefinition", err)) + } + return ctrl.Result{}, nil + } + + if err = r.createOrUpdateWFStepDefRevision(ctx, req.Namespace, &wfstepdefinition, defRev); err != nil { + klog.ErrorS(err, "cannot create DefinitionRevision") + r.record.Event(&(wfstepdefinition), event.Warning("cannot create DefinitionRevision", err)) + return ctrl.Result{}, util.PatchCondition(ctx, r, &(wfstepdefinition), + cpv1alpha1.ReconcileError(fmt.Errorf(util.ErrCreateOrUpdateDefinitionRevision, defRev.Name, err))) + } + klog.InfoS("Successfully createOrUpdateWFStepDefRevision", "name", defRev.Name) + + wfstepdefinition.Status.LatestRevision = &common.Revision{ + Name: defRev.Name, + Revision: defRev.Spec.Revision, + RevisionHash: defRev.Spec.RevisionHash, + } + + if err := r.UpdateStatus(ctx, &wfstepdefinition); err != nil { + klog.ErrorS(err, "cannot update WorkflowStepDefinition Status") + r.record.Event(&(wfstepdefinition), event.Warning("cannot update WorkflowStepDefinition Status", err)) + return ctrl.Result{}, util.PatchCondition(ctx, r, &(wfstepdefinition), + cpv1alpha1.ReconcileError(fmt.Errorf(util.ErrUpdateWorkflowStepDefinition, wfstepdefinition.Name, err))) + } + + if err := coredef.CleanUpDefinitionRevision(ctx, r.Client, &wfstepdefinition, r.defRevLimit); err != nil { + klog.Error("[Garbage collection]") + r.record.Event(&wfstepdefinition, event.Warning("failed to garbage collect DefinitionRevision of type WorkflowStepDefinition", err)) + } + + return ctrl.Result{}, nil +} + +func (r *Reconciler) createOrUpdateWFStepDefRevision(ctx context.Context, ns string, + def *v1beta1.WorkflowStepDefinition, defRev *v1beta1.DefinitionRevision) error { + + ownerReference := []metav1.OwnerReference{{ + APIVersion: def.APIVersion, + Kind: def.Kind, + Name: def.Name, + UID: def.GetUID(), + Controller: pointer.BoolPtr(true), + BlockOwnerDeletion: pointer.BoolPtr(true), + }} + + defRev.SetLabels(def.GetLabels()) + defRev.SetLabels(util.MergeMapOverrideWithDst(defRev.Labels, + map[string]string{oam.LabelWorkflowStepDefinitionName: def.Name})) + defRev.SetNamespace(ns) + defRev.SetAnnotations(def.GetAnnotations()) + defRev.SetOwnerReferences(ownerReference) + + rev := &v1beta1.DefinitionRevision{} + if err := r.Client.Get(ctx, client.ObjectKey{Namespace: ns, Name: defRev.Name}, rev); err != nil { + if apierrors.IsNotFound(err) { + return r.Create(ctx, defRev) + } + return err + } + + rev.SetAnnotations(defRev.GetAnnotations()) + rev.SetLabels(defRev.GetLabels()) + rev.SetOwnerReferences(ownerReference) + return r.Update(ctx, rev) +} + +// UpdateStatus updates v1beta1.WorkflowStepDefinition's Status with retry.RetryOnConflict +func (r *Reconciler) UpdateStatus(ctx context.Context, def *v1beta1.WorkflowStepDefinition, opts ...client.UpdateOption) error { + status := def.DeepCopy().Status + return retry.RetryOnConflict(retry.DefaultBackoff, func() (err error) { + if err = r.Get(ctx, client.ObjectKey{Namespace: def.Namespace, Name: def.Name}, def); err != nil { + return + } + def.Status = status + return r.Status().Update(ctx, def, opts...) + }) +} + +// SetupWithManager will setup with event recorder +func (r *Reconciler) SetupWithManager(mgr ctrl.Manager) error { + r.record = event.NewAPIRecorder(mgr.GetEventRecorderFor("WorkflowStepDefinition")). + WithAnnotations("controller", "WorkflowStepDefinition") + return ctrl.NewControllerManagedBy(mgr). + For(&v1beta1.WorkflowStepDefinition{}). + Complete(r) +} + +// Setup adds a controller that reconciles WorkflowStepDefinition. +func Setup(mgr ctrl.Manager, args controller.Args, _ logging.Logger) error { + r := Reconciler{ + Client: mgr.GetClient(), + Scheme: mgr.GetScheme(), + dm: args.DiscoveryMapper, + pd: args.PackageDiscover, + defRevLimit: args.DefRevisionLimit, + } + return r.SetupWithManager(mgr) +} diff --git a/pkg/controller/core.oam.dev/v1alpha2/setup.go b/pkg/controller/core.oam.dev/v1alpha2/setup.go index d0926c41f..cd66a3779 100644 --- a/pkg/controller/core.oam.dev/v1alpha2/setup.go +++ b/pkg/controller/core.oam.dev/v1alpha2/setup.go @@ -27,9 +27,11 @@ import ( "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/applicationcontext" "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/applicationrollout" "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/core/components/componentdefinition" + "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/core/policies/policydefinition" "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/core/scopes/healthscope" "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/core/traits/manualscalertrait" "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/core/traits/traitdefinition" + "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/core/workflow/workflowstepdefinition" "github.com/oam-dev/kubevela/pkg/controller/core.oam.dev/v1alpha2/core/workloads/containerizedworkload" ) @@ -38,7 +40,7 @@ func Setup(mgr ctrl.Manager, args controller.Args, l logging.Logger) error { for _, setup := range []func(ctrl.Manager, controller.Args, logging.Logger) error{ containerizedworkload.Setup, manualscalertrait.Setup, healthscope.Setup, application.Setup, applicationrollout.Setup, applicationcontext.Setup, appdeployment.Setup, - traitdefinition.Setup, componentdefinition.Setup, + traitdefinition.Setup, componentdefinition.Setup, policydefinition.Setup, workflowstepdefinition.Setup, } { if err := setup(mgr, args, l); err != nil { return err diff --git a/pkg/oam/labels.go b/pkg/oam/labels.go index 8059aea18..6a5e21a4a 100644 --- a/pkg/oam/labels.go +++ b/pkg/oam/labels.go @@ -43,7 +43,7 @@ const ( // LabelComponentDefinitionName records the name of ComponentDefinition LabelComponentDefinitionName = "componentdefinition.oam.dev/name" - // LabelComponentDefinitionName records the name of TraitDefinition + // LabelTraitDefinitionName records the name of TraitDefinition LabelTraitDefinitionName = "trait.oam.dev/name" // LabelPolicyDefinitionName records the name of PolicyDefinition LabelPolicyDefinitionName = "policydefinition.oam.dev/name"