diff --git a/pkg/apiserver/rest/usecase/application.go b/pkg/apiserver/rest/usecase/application.go index eef22d55c..1ed2b8dea 100644 --- a/pkg/apiserver/rest/usecase/application.go +++ b/pkg/apiserver/rest/usecase/application.go @@ -33,18 +33,18 @@ import ( "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/selection" "k8s.io/apimachinery/pkg/types" + "k8s.io/klog/v2" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/yaml" "github.com/oam-dev/kubevela/apis/core.oam.dev/common" "github.com/oam-dev/kubevela/apis/core.oam.dev/v1alpha1" "github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1" + velatypes "github.com/oam-dev/kubevela/apis/types" "github.com/oam-dev/kubevela/pkg/apiserver/clients" "github.com/oam-dev/kubevela/pkg/apiserver/datastore" "github.com/oam-dev/kubevela/pkg/apiserver/log" "github.com/oam-dev/kubevela/pkg/apiserver/model" - - velatypes "github.com/oam-dev/kubevela/apis/types" apisv1 "github.com/oam-dev/kubevela/pkg/apiserver/rest/apis/v1" "github.com/oam-dev/kubevela/pkg/apiserver/rest/utils" "github.com/oam-dev/kubevela/pkg/apiserver/rest/utils/bcode" @@ -705,21 +705,9 @@ func (c *applicationUsecaseImpl) Deploy(ctx context.Context, app *model.Applicat } // sync configs to clusters - // TODO(zzxwill) need to check the type of the componentDefinition, if it is `Cloud`, skip the sync - targets, err := listTarget(ctx, c.ds, app.Project, nil) - if err != nil { + if err := c.syncConfigs4Application(ctx, oamApp, app.Project, workflow.EnvName); err != nil { return nil, err } - var clusterTargets []*model.ClusterTarget - for i, t := range targets { - if t.Cluster != nil { - clusterTargets = append(clusterTargets, targets[i].Cluster) - } - } - - if err := SyncConfigs(ctx, c.kubeClient, app.Project, clusterTargets); err != nil { - return nil, fmt.Errorf("sync config failure %w", err) - } // step2: check and create deploy event if !req.Force { @@ -809,6 +797,44 @@ func (c *applicationUsecaseImpl) Deploy(ctx context.Context, app *model.Applicat }, nil } +// sync configs to clusters +func (c *applicationUsecaseImpl) syncConfigs4Application(ctx context.Context, app *v1beta1.Application, projectName, envName string) error { + var areTerraformComponents = true + for _, m := range app.Spec.Components { + d := &v1beta1.ComponentDefinition{} + if err := c.kubeClient.Get(ctx, client.ObjectKey{Namespace: velatypes.DefaultKubeVelaNS, Name: m.Type}, d); err != nil { + klog.ErrorS(err, "failed to get config type", "ComponentDefinition", m.Type) + } + // check the type of the componentDefinition is Terraform + if d.Spec.Schematic != nil && d.Spec.Schematic.Terraform == nil { + areTerraformComponents = false + } + } + // skip configs sync + if areTerraformComponents { + return nil + } + env, err := c.envUsecase.GetEnv(ctx, envName) + if err != nil { + return err + } + var clusterTargets []*model.ClusterTarget + for _, t := range env.Targets { + target, err := c.targetUsecase.GetTarget(ctx, t) + if err != nil { + return err + } + if target.Cluster != nil { + clusterTargets = append(clusterTargets, target.Cluster) + } + } + + if err := SyncConfigs(ctx, c.kubeClient, projectName, clusterTargets); err != nil { + return fmt.Errorf("sync config failure %w", err) + } + return nil +} + func (c *applicationUsecaseImpl) renderOAMApplication(ctx context.Context, appModel *model.Application, reqWorkflowName, version string) (*v1beta1.Application, error) { // Priority 1 uses the requested workflow as release . // Priority 2 uses the default workflow as release . diff --git a/pkg/apiserver/rest/usecase/config.go b/pkg/apiserver/rest/usecase/config.go index 9be91a4cd..4fac6f7ec 100644 --- a/pkg/apiserver/rest/usecase/config.go +++ b/pkg/apiserver/rest/usecase/config.go @@ -21,11 +21,13 @@ import ( "encoding/json" "fmt" + set "github.com/deckarep/golang-set" "github.com/pkg/errors" v1 "k8s.io/api/core/v1" kerrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" + "k8s.io/klog/v2" "sigs.k8s.io/controller-runtime/pkg/client" "github.com/oam-dev/kubevela/apis/core.oam.dev/common" @@ -41,10 +43,11 @@ const ( definitionAlias = definition.UserPrefix + "alias.config.oam.dev" definitionType = definition.UserPrefix + "type.config.oam.dev" - velaCoreConfig = "velacore-config" - configIsReady = "Ready" - configIsNotReady = "Not ready" - terraformProviderAlias = "Terraform Cloud Provider" + velaCoreConfig = "velacore-config" + configIsReady = "Ready" + configIsNotReady = "Not ready" + terraformProviderAlias = "Terraform Cloud Provider" + configSyncProjectPrefix = "config-sync" ) // ConfigHandler handle CRUD of configs @@ -246,13 +249,12 @@ type ApplicationDeployTarget struct { // SyncConfigs will sync configs to working clusters func SyncConfigs(ctx context.Context, k8sClient client.Client, project string, targets []*model.ClusterTarget) error { - name := fmt.Sprintf("config-sync-%s", project) + name := fmt.Sprintf("%s-%s", configSyncProjectPrefix, project) // get all configs which can be synced to working clusters in the project var secrets v1.SecretList if err := k8sClient.List(ctx, &secrets, client.InNamespace(types.DefaultKubeVelaNS), client.MatchingLabels{ types.LabelConfigCatalog: velaCoreConfig, - types.LabelConfigProject: project, types.LabelConfigSyncToMultiCluster: "true", }); err != nil { return err @@ -262,9 +264,11 @@ func SyncConfigs(ctx context.Context, k8sClient client.Client, project string, t } objects := make([]map[string]string, len(secrets.Items)) for i, s := range secrets.Items { - objects[i] = map[string]string{ - "name": s.Name, - "resource": "secret", + if s.Labels[types.LabelConfigProject] == "" || s.Labels[types.LabelConfigProject] == project { + objects[i] = map[string]string{ + "name": s.Name, + "resource": "secret", + } } } objectsBytes, err := json.Marshal(map[string][]map[string]string{"objects": objects}) @@ -278,6 +282,26 @@ func SyncConfigs(ctx context.Context, k8sClient client.Client, project string, t return err } // config sync application doesn't exist, create one + clusterTargets := convertClusterTargets(targets) + if len(clusterTargets) == 0 { + errMsg := "no policy (no targets found) to sync configs" + klog.InfoS(errMsg, "project", project) + return errors.New(errMsg) + } + policies := make([]v1beta1.AppPolicy, len(clusterTargets)) + for i, t := range clusterTargets { + properties, err := json.Marshal(t) + if err != nil { + return err + } + policies[i] = v1beta1.AppPolicy{ + Type: "topology", + Name: t.Namespace, + Properties: &runtime.RawExtension{ + Raw: properties, + }, + } + } scratch := &v1beta1.Application{ ObjectMeta: metav1.ObjectMeta{ @@ -297,11 +321,10 @@ func SyncConfigs(ctx context.Context, k8sClient client.Client, project string, t Properties: &runtime.RawExtension{Raw: objectsBytes}, }, }, + Policies: policies, }, } - if err := k8sClient.Create(ctx, scratch); err != nil { - return err - } + return k8sClient.Create(ctx, scratch) } // config sync application exists, update it app.Spec.Components = []common.ApplicationComponent{ @@ -321,6 +344,11 @@ func SyncConfigs(ctx context.Context, k8sClient client.Client, project string, t } mergedTarget := mergeTargets(currentTargets, targets) + if len(mergedTarget) == 0 { + errMsg := "no policy (no targets found) to sync configs" + klog.InfoS(errMsg, "project", project) + return errors.New(errMsg) + } mergedPolicies := make([]v1beta1.AppPolicy, len(mergedTarget)) for i, t := range mergedTarget { properties, err := json.Marshal(t) @@ -340,14 +368,23 @@ func SyncConfigs(ctx context.Context, k8sClient client.Client, project string, t } func mergeTargets(currentTargets []ApplicationDeployTarget, targets []*model.ClusterTarget) []ApplicationDeployTarget { - var mergedTargets []ApplicationDeployTarget + var ( + mergedTargets []ApplicationDeployTarget + // make sure the clusters of target with same namespace are merged + clusterTargets = convertClusterTargets(targets) + ) + for _, c := range currentTargets { var hasSameNamespace bool - for _, t := range targets { + for _, t := range clusterTargets { if c.Namespace == t.Namespace { hasSameNamespace = true - clusters := append(c.Clusters, t.ClusterName) - mergedTargets = append(mergedTargets, ApplicationDeployTarget{Namespace: c.Namespace, Clusters: clusters}) + clusters := set.NewSetFromSlice(stringToInterfaceSlice(t.Clusters)) + for _, cluster := range c.Clusters { + clusters.Add(cluster) + } + mergedTargets = append(mergedTargets, ApplicationDeployTarget{Namespace: c.Namespace, + Clusters: interfaceToStringSlice(clusters.ToSlice())}) } } if !hasSameNamespace { @@ -355,7 +392,7 @@ func mergeTargets(currentTargets []ApplicationDeployTarget, targets []*model.Clu } } - for _, t := range targets { + for _, t := range clusterTargets { var hasSameNamespace bool for _, c := range currentTargets { if c.Namespace == t.Namespace { @@ -363,9 +400,75 @@ func mergeTargets(currentTargets []ApplicationDeployTarget, targets []*model.Clu } } if !hasSameNamespace { - mergedTargets = append(mergedTargets, ApplicationDeployTarget{Namespace: t.Namespace, Clusters: []string{t.ClusterName}}) + mergedTargets = append(mergedTargets, t) } } return mergedTargets } + +func convertClusterTargets(targets []*model.ClusterTarget) []ApplicationDeployTarget { + type Target struct { + Namespace string `json:"namespace"` + Clusters []interface{} `json:"clusters"` + } + + var ( + clusterTargets []Target + namespaceSet = set.NewSet() + ) + + for i := 0; i < len(targets); i++ { + clusters := set.NewSet(targets[i].ClusterName) + for j := i + 1; j < len(targets); j++ { + if targets[i].Namespace == targets[j].Namespace { + clusters.Add(targets[j].ClusterName) + } + } + if namespaceSet.Contains(targets[i].Namespace) { + continue + } + clusterTargets = append(clusterTargets, Target{ + Namespace: targets[i].Namespace, + Clusters: clusters.ToSlice(), + }) + namespaceSet.Add(targets[i].Namespace) + } + + t := make([]ApplicationDeployTarget, len(clusterTargets)) + for i, ct := range clusterTargets { + t[i] = ApplicationDeployTarget{ + Namespace: ct.Namespace, + Clusters: interfaceToStringSlice(ct.Clusters), + } + } + return t +} + +func interfaceToStringSlice(i []interface{}) []string { + var s []string + for _, v := range i { + s = append(s, v.(string)) + } + return s +} + +func stringToInterfaceSlice(i []string) []interface{} { + var s []interface{} + for _, v := range i { + s = append(s, v) + } + return s +} + +// destroySyncConfigsApp will delete the application which is used to sync configs +func destroySyncConfigsApp(ctx context.Context, k8sClient client.Client, project string) error { + name := fmt.Sprintf("%s-%s", configSyncProjectPrefix, project) + var app = &v1beta1.Application{} + if err := k8sClient.Get(ctx, client.ObjectKey{Namespace: types.DefaultKubeVelaNS, Name: name}, app); err != nil { + if !kerrors.IsNotFound(err) { + return err + } + } + return k8sClient.Delete(ctx, app) +} diff --git a/pkg/apiserver/rest/usecase/config_test.go b/pkg/apiserver/rest/usecase/config_test.go index f8efded4f..7fbd18f7b 100644 --- a/pkg/apiserver/rest/usecase/config_test.go +++ b/pkg/apiserver/rest/usecase/config_test.go @@ -18,6 +18,9 @@ package usecase import ( "context" + "encoding/json" + "sort" + "strings" "testing" . "github.com/agiledragon/gomonkey/v2" @@ -31,6 +34,7 @@ import ( "github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1" "github.com/oam-dev/kubevela/apis/types" + "github.com/oam-dev/kubevela/pkg/apiserver/model" apis "github.com/oam-dev/kubevela/pkg/apiserver/rest/apis/v1" "github.com/oam-dev/kubevela/pkg/definition" "github.com/oam-dev/kubevela/pkg/multicluster" @@ -337,3 +341,239 @@ func TestGetConfigs(t *testing.T) { }) } } + +func TestMergeTargets(t *testing.T) { + currentTargets := []ApplicationDeployTarget{ + { + Namespace: "n1", + Clusters: []string{"c1", "c2"}, + }, { + Namespace: "n2", + Clusters: []string{"c3"}, + }, + } + targets := []*model.ClusterTarget{ + { + Namespace: "n3", + ClusterName: "c4", + }, { + Namespace: "n1", + ClusterName: "c5", + }, + { + Namespace: "n2", + ClusterName: "c3", + }, + } + + expected := []ApplicationDeployTarget{ + { + Namespace: "n1", + Clusters: []string{"c1", "c2", "c5"}, + }, { + Namespace: "n2", + Clusters: []string{"c3"}, + }, { + Namespace: "n3", + Clusters: []string{"c4"}, + }, + } + + got := mergeTargets(currentTargets, targets) + + for i, g := range got { + clusters := g.Clusters + sort.SliceStable(clusters, func(i, j int) bool { + return clusters[i] < clusters[j] + }) + got[i].Clusters = clusters + } + assert.DeepEqual(t, expected, got) +} + +func TestConvert(t *testing.T) { + targets := []*model.ClusterTarget{ + { + Namespace: "n3", + ClusterName: "c4", + }, { + Namespace: "n1", + ClusterName: "c5", + }, + { + Namespace: "n2", + ClusterName: "c3", + }, + { + Namespace: "n3", + ClusterName: "c5", + }, + } + + expected := []ApplicationDeployTarget{ + { + Namespace: "n3", + Clusters: []string{"c4", "c5"}, + }, + { + Namespace: "n1", + Clusters: []string{"c5"}, + }, { + Namespace: "n2", + Clusters: []string{"c3"}, + }, + } + + got := convertClusterTargets(targets) + + for i, g := range got { + clusters := g.Clusters + sort.SliceStable(clusters, func(i, j int) bool { + return clusters[i] < clusters[j] + }) + got[i].Clusters = clusters + } + assert.DeepEqual(t, expected, got) +} + +func TestDestroySyncConfigsApp(t *testing.T) { + s := runtime.NewScheme() + v1beta1.AddToScheme(s) + corev1.AddToScheme(s) + app1 := &v1beta1.Application{ + ObjectMeta: metav1.ObjectMeta{ + Name: "config-sync-p1", + Namespace: types.DefaultKubeVelaNS, + }, + } + k8sClient1 := fake.NewClientBuilder().WithScheme(s).WithObjects(app1).Build() + + k8sClient2 := fake.NewClientBuilder().Build() + + type args struct { + project string + k8sClient client.Client + } + + type want struct { + errMsg string + } + + ctx := context.Background() + + testcases := map[string]struct { + args args + want want + }{ + "found": { + args: args{ + project: "p1", + k8sClient: k8sClient1, + }, + }, + "not found": { + args: args{ + project: "p1", + k8sClient: k8sClient2, + }, + want: want{ + errMsg: "no kind is registered for the type v1beta1.Application", + }, + }, + } + for name, tc := range testcases { + t.Run(name, func(t *testing.T) { + err := destroySyncConfigsApp(ctx, tc.args.k8sClient, tc.args.project) + if err != nil || tc.want.errMsg != "" { + if !strings.Contains(err.Error(), tc.want.errMsg) { + assert.ErrorContains(t, err, tc.want.errMsg) + } + } + }) + } +} + +func TestSyncConfigs(t *testing.T) { + s := runtime.NewScheme() + v1beta1.AddToScheme(s) + corev1.AddToScheme(s) + secret1 := &corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{ + Name: "s1", + Namespace: types.DefaultKubeVelaNS, + Labels: map[string]string{ + types.LabelConfigCatalog: velaCoreConfig, + types.LabelConfigProject: "p1", + types.LabelConfigSyncToMultiCluster: "true", + }, + }, + } + + policies := []ApplicationDeployTarget{{ + Namespace: "n9", + Clusters: []string{"c19"}, + }} + properties, _ := json.Marshal(policies) + app1 := &v1beta1.Application{ + ObjectMeta: metav1.ObjectMeta{ + Name: "config-sync-p2", + Namespace: types.DefaultKubeVelaNS, + }, + Spec: v1beta1.ApplicationSpec{ + Policies: []v1beta1.AppPolicy{{ + Name: "c19", + Type: "topology", + Properties: &runtime.RawExtension{Raw: properties}, + }}, + }, + } + + k8sClient := fake.NewClientBuilder().WithScheme(s).WithObjects(secret1, app1).Build() + + type args struct { + project string + targets []*model.ClusterTarget + } + + type want struct { + errMsg string + } + + ctx := context.Background() + + testcases := []struct { + name string + args args + want want + }{ + { + name: "create", + args: args{ + project: "p1", + targets: []*model.ClusterTarget{{ + ClusterName: "c1", + Namespace: "n1", + }}, + }, + }, + { + name: "update", + args: args{ + project: "p2", + targets: []*model.ClusterTarget{{ + ClusterName: "c1", + Namespace: "n1", + }}, + }, + }, + } + + for _, tc := range testcases { + t.Run(tc.name, func(t *testing.T) { + err := SyncConfigs(ctx, k8sClient, tc.args.project, tc.args.targets) + if tc.want.errMsg != "" || err != nil { + assert.ErrorContains(t, err, tc.want.errMsg) + } + }) + } +} diff --git a/pkg/apiserver/rest/usecase/project.go b/pkg/apiserver/rest/usecase/project.go index 206b5d5f9..46f3f3c42 100644 --- a/pkg/apiserver/rest/usecase/project.go +++ b/pkg/apiserver/rest/usecase/project.go @@ -298,7 +298,11 @@ func (p *projectUsecaseImpl) DeleteProject(ctx context.Context, name string) err return err } } - return p.ds.Delete(ctx, &model.Project{Name: name}) + if err := p.ds.Delete(ctx, &model.Project{Name: name}); err != nil { + return err + } + // delete config-sync application + return destroySyncConfigsApp(ctx, p.k8sClient, name) } // CreateProject create project diff --git a/pkg/apiserver/rest/usecase/project_test.go b/pkg/apiserver/rest/usecase/project_test.go index 197ec838b..e714b76d5 100644 --- a/pkg/apiserver/rest/usecase/project_test.go +++ b/pkg/apiserver/rest/usecase/project_test.go @@ -164,6 +164,18 @@ var _ = Describe("Test project usecase functions", func() { Name: "test-project", Description: "this is a project description", } + app1 := &v1beta1.Application{ + ObjectMeta: metav1.ObjectMeta{ + Name: "config-sync-test-project", + Namespace: "vela-system", + }, + Spec: v1beta1.ApplicationSpec{ + Components: []common.ApplicationComponent{{ + Type: "aaa", + }}, + }, + } + Expect(k8sClient.Create(context.TODO(), app1)).Should(BeNil()) _, err := projectUsecase.CreateProject(context.TODO(), req) Expect(err).Should(BeNil()) @@ -232,6 +244,19 @@ var _ = Describe("Test project usecase functions", func() { Name: "test-project", Description: "this is a project description", } + app1 := &v1beta1.Application{ + ObjectMeta: metav1.ObjectMeta{ + Name: "config-sync-test-project", + Namespace: "vela-system", + }, + Spec: v1beta1.ApplicationSpec{ + Components: []common.ApplicationComponent{{ + Type: "aaa", + }}, + }, + } + Expect(k8sClient.Create(context.TODO(), app1)).Should(BeNil()) + _, err := projectUsecase.CreateProject(context.TODO(), req) Expect(err).Should(BeNil())