mirror of
https://github.com/kubevela/kubevela.git
synced 2026-08-27 16:17:34 +00:00
add PolicyDefinition and WorkflowStepDefinition controller (#1751)
This commit is contained in:
+295
@@ -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{},
|
||||
},
|
||||
},
|
||||
}
|
||||
+205
@@ -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)
|
||||
}
|
||||
@@ -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())
|
||||
}
|
||||
+295
@@ -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{},
|
||||
},
|
||||
},
|
||||
}
|
||||
+121
@@ -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())
|
||||
}
|
||||
+205
@@ -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)
|
||||
}
|
||||
@@ -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
|
||||
|
||||
+1
-1
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user