diff --git a/pkg/app/piped/cloudprovider/kubernetes/BUILD.bazel b/pkg/app/piped/cloudprovider/kubernetes/BUILD.bazel index 6fa65d3729..3df2ce4070 100644 --- a/pkg/app/piped/cloudprovider/kubernetes/BUILD.bazel +++ b/pkg/app/piped/cloudprovider/kubernetes/BUILD.bazel @@ -3,6 +3,7 @@ load("@io_bazel_rules_go//go:def.bzl", "go_library", "go_test") go_library( name = "go_default_library", srcs = [ + "applier.go", "cache.go", "deployment.go", "diff.go", @@ -12,6 +13,7 @@ go_library( "kubectl.go", "kubernetes.go", "kustomize.go", + "loader.go", "manifest.go", "resourcekey.go", "state.go", diff --git a/pkg/app/piped/cloudprovider/kubernetes/applier.go b/pkg/app/piped/cloudprovider/kubernetes/applier.go new file mode 100644 index 0000000000..b4a37df32c --- /dev/null +++ b/pkg/app/piped/cloudprovider/kubernetes/applier.go @@ -0,0 +1,131 @@ +// Copyright 2022 The PipeCD 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 kubernetes + +import ( + "context" + "errors" + "fmt" + "sync" + + "go.uber.org/zap" + + "github.com/pipe-cd/pipecd/pkg/app/piped/toolregistry" + "github.com/pipe-cd/pipecd/pkg/config" +) + +type Applier interface { + // ApplyManifest does applying the given manifest. + ApplyManifest(ctx context.Context, manifest Manifest) error + // CreateManifest does creating resource from given manifest. + CreateManifest(ctx context.Context, manifest Manifest) error + // ReplaceManifest does replacing resource from given manifest. + ReplaceManifest(ctx context.Context, manifest Manifest) error + // Delete deletes the given resource from Kubernetes cluster. + Delete(ctx context.Context, key ResourceKey) error +} + +type applier struct { + input config.KubernetesDeploymentInput + logger *zap.Logger + + kubectl *Kubectl + initOnce sync.Once + initErr error +} + +func NewApplier(input config.KubernetesDeploymentInput, logger *zap.Logger) Applier { + return &applier{ + input: input, + logger: logger.Named("kubernetes-applier"), + } +} + +// ApplyManifest does applying the given manifest. +func (a *applier) ApplyManifest(ctx context.Context, manifest Manifest) error { + a.initOnce.Do(func() { + a.kubectl, a.initErr = a.findKubectl(ctx, a.input.KubectlVersion) + }) + if a.initErr != nil { + return a.initErr + } + + return a.kubectl.Apply(ctx, a.getNamespaceToRun(manifest.Key), manifest) +} + +// CreateManifest uses kubectl to create the given manifests. +func (a *applier) CreateManifest(ctx context.Context, manifest Manifest) error { + a.initOnce.Do(func() { + a.kubectl, a.initErr = a.findKubectl(ctx, a.input.KubectlVersion) + }) + if a.initErr != nil { + return a.initErr + } + + return a.kubectl.Create(ctx, a.getNamespaceToRun(manifest.Key), manifest) +} + +// ReplaceManifest uses kubectl to replace the given manifests. +func (a *applier) ReplaceManifest(ctx context.Context, manifest Manifest) error { + a.initOnce.Do(func() { + a.kubectl, a.initErr = a.findKubectl(ctx, a.input.KubectlVersion) + }) + if a.initErr != nil { + return a.initErr + } + + err := a.kubectl.Replace(ctx, a.getNamespaceToRun(manifest.Key), manifest) + if err == nil { + return nil + } + + if errors.Is(err, errorReplaceNotFound) { + return ErrNotFound + } + + return err +} + +// Delete deletes the given resource from Kubernetes cluster. +func (a *applier) Delete(ctx context.Context, k ResourceKey) (err error) { + a.initOnce.Do(func() { + a.kubectl, a.initErr = a.findKubectl(ctx, a.input.KubectlVersion) + }) + if a.initErr != nil { + return a.initErr + } + + return a.kubectl.Delete(ctx, a.getNamespaceToRun(k), k) +} + +// getNamespaceToRun returns namespace used on kubectl apply/delete commands. +// priority: config.KubernetesDeploymentInput > kubernetes.ResourceKey +func (a *applier) getNamespaceToRun(k ResourceKey) string { + if a.input.Namespace != "" { + return a.input.Namespace + } + return k.Namespace +} + +func (a *applier) findKubectl(ctx context.Context, version string) (*Kubectl, error) { + path, installed, err := toolregistry.DefaultRegistry().Kubectl(ctx, version) + if err != nil { + return nil, fmt.Errorf("no kubectl %s (%v)", version, err) + } + if installed { + a.logger.Info(fmt.Sprintf("kubectl %s has just been installed because of no pre-installed binary for that version", version)) + } + return NewKubectl(version, path), nil +} diff --git a/pkg/app/piped/cloudprovider/kubernetes/kubernetes.go b/pkg/app/piped/cloudprovider/kubernetes/kubernetes.go index 4a1770db7d..454c25eeec 100644 --- a/pkg/app/piped/cloudprovider/kubernetes/kubernetes.go +++ b/pkg/app/piped/cloudprovider/kubernetes/kubernetes.go @@ -1,4 +1,4 @@ -// Copyright 2020 The PipeCD Authors. +// Copyright 2022 The PipeCD Authors. // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. @@ -15,18 +15,7 @@ package kubernetes import ( - "context" "errors" - "fmt" - "os" - "path/filepath" - "sync" - - "go.uber.org/zap" - - "github.com/pipe-cd/pipecd/pkg/app/piped/toolregistry" - "github.com/pipe-cd/pipecd/pkg/config" - "github.com/pipe-cd/pipecd/pkg/git" ) var ( @@ -49,279 +38,3 @@ const ( kustomizationFileName = "kustomization.yaml" ) - -type TemplatingMethod string - -const ( - TemplatingMethodHelm TemplatingMethod = "helm" - TemplatingMethodKustomize TemplatingMethod = "kustomize" - TemplatingMethodNone TemplatingMethod = "none" -) - -type Provider interface { - ManifestLoader - Applier -} - -type ManifestLoader interface { - // LoadManifests renders and loads all manifests for application. - LoadManifests(ctx context.Context) ([]Manifest, error) -} - -type Applier interface { - // Apply does applying application manifests by using the tool specified in Input. - Apply(ctx context.Context) error - // ApplyManifest does applying the given manifest. - ApplyManifest(ctx context.Context, manifest Manifest) error - // CreateManifest does creating resource from given manifest. - CreateManifest(ctx context.Context, manifest Manifest) error - // ReplaceManifest does replacing resource from given manifest. - ReplaceManifest(ctx context.Context, manifest Manifest) error - // Delete deletes the given resource from Kubernetes cluster. - Delete(ctx context.Context, key ResourceKey) error -} - -type gitClient interface { - Clone(ctx context.Context, repoID, remote, branch, destination string) (git.Repo, error) -} - -type provider struct { - appName string - appDir string - repoDir string - configFileName string - input config.KubernetesDeploymentInput - gc gitClient - logger *zap.Logger - - kubectl *Kubectl - kustomize *Kustomize - helm *Helm - templatingMethod TemplatingMethod - initOnce sync.Once - initErr error -} - -func NewProvider( - appName, appDir, repoDir, configFileName string, - input config.KubernetesDeploymentInput, - gc gitClient, - logger *zap.Logger, -) Provider { - - return &provider{ - appName: appName, - appDir: appDir, - repoDir: repoDir, - configFileName: configFileName, - input: input, - gc: gc, - logger: logger.Named("kubernetes-provider"), - } -} - -func NewManifestLoader( - appName, appDir, repoDir, configFileName string, - input config.KubernetesDeploymentInput, - gc gitClient, - logger *zap.Logger, -) ManifestLoader { - - return NewProvider(appName, appDir, repoDir, configFileName, input, gc, logger) -} - -func (p *provider) init(ctx context.Context) { - p.templatingMethod = determineTemplatingMethod(p.input, p.appDir) - - // We need kubectl for all templating methods. - p.kubectl, p.initErr = p.findKubectl(ctx, p.input.KubectlVersion) - if p.initErr != nil { - return - } - - switch p.templatingMethod { - case TemplatingMethodHelm: - p.helm, p.initErr = p.findHelm(ctx, p.input.HelmVersion) - - case TemplatingMethodKustomize: - p.kustomize, p.initErr = p.findKustomize(ctx, p.input.KustomizeVersion) - } -} - -// LoadManifests renders and loads all manifests for application. -func (p *provider) LoadManifests(ctx context.Context) (manifests []Manifest, err error) { - p.initOnce.Do(func() { p.init(ctx) }) - if p.initErr != nil { - return nil, p.initErr - } - - switch p.templatingMethod { - case TemplatingMethodHelm: - var data string - switch { - case p.input.HelmChart.GitRemote != "": - chart := helmRemoteGitChart{ - GitRemote: p.input.HelmChart.GitRemote, - Ref: p.input.HelmChart.Ref, - Path: p.input.HelmChart.Path, - } - data, err = p.helm.TemplateRemoteGitChart(ctx, - p.appName, - p.appDir, - p.input.Namespace, - chart, - p.gc, - p.input.HelmOptions) - - case p.input.HelmChart.Repository != "": - chart := helmRemoteChart{ - Repository: p.input.HelmChart.Repository, - Name: p.input.HelmChart.Name, - Version: p.input.HelmChart.Version, - Insecure: p.input.HelmChart.Insecure, - } - data, err = p.helm.TemplateRemoteChart(ctx, - p.appName, - p.appDir, - p.input.Namespace, - chart, - p.input.HelmOptions) - - default: - data, err = p.helm.TemplateLocalChart(ctx, - p.appName, - p.appDir, - p.input.Namespace, - p.input.HelmChart.Path, - p.input.HelmOptions) - } - - if err != nil { - err = fmt.Errorf("unable to run helm template: %w", err) - return - } - manifests, err = ParseManifests(data) - - case TemplatingMethodKustomize: - var data string - data, err = p.kustomize.Template(ctx, p.appName, p.appDir, p.input.KustomizeOptions) - if err != nil { - err = fmt.Errorf("unable to run kustomize template: %w", err) - return - } - manifests, err = ParseManifests(data) - - case TemplatingMethodNone: - manifests, err = LoadPlainYAMLManifests(p.appDir, p.input.Manifests, p.configFileName) - - default: - err = fmt.Errorf("unsupport templating method %v", p.templatingMethod) - } - - return -} - -// Apply does applying application manifests by using the tool specified in Input. -func (p *provider) Apply(ctx context.Context) error { - return nil -} - -// ApplyManifest does applying the given manifest. -func (p *provider) ApplyManifest(ctx context.Context, manifest Manifest) error { - p.initOnce.Do(func() { p.init(ctx) }) - if p.initErr != nil { - return p.initErr - } - - return p.kubectl.Apply(ctx, p.getNamespaceToRun(manifest.Key), manifest) -} - -// CreateManifest uses kubectl to create the given manifests. -func (p *provider) CreateManifest(ctx context.Context, manifest Manifest) error { - p.initOnce.Do(func() { p.init(ctx) }) - if p.initErr != nil { - return p.initErr - } - - return p.kubectl.Create(ctx, p.getNamespaceToRun(manifest.Key), manifest) -} - -// ReplaceManifest uses kubectl to replace the given manifests. -func (p *provider) ReplaceManifest(ctx context.Context, manifest Manifest) error { - p.initOnce.Do(func() { p.init(ctx) }) - if p.initErr != nil { - return p.initErr - } - err := p.kubectl.Replace(ctx, p.getNamespaceToRun(manifest.Key), manifest) - if err == nil { - return nil - } - - if errors.Is(err, errorReplaceNotFound) { - return ErrNotFound - } - - return err -} - -// Delete deletes the given resource from Kubernetes cluster. -func (p *provider) Delete(ctx context.Context, k ResourceKey) (err error) { - p.initOnce.Do(func() { p.init(ctx) }) - if p.initErr != nil { - return p.initErr - } - - return p.kubectl.Delete(ctx, p.getNamespaceToRun(k), k) -} - -// getNamespaceToRun returns namespace used on kubectl apply/delete commands. -// priority: config.KubernetesDeploymentInput > kubernetes.ResourceKey -func (p *provider) getNamespaceToRun(k ResourceKey) string { - if p.input.Namespace != "" { - return p.input.Namespace - } - return k.Namespace -} - -func (p *provider) findKubectl(ctx context.Context, version string) (*Kubectl, error) { - path, installed, err := toolregistry.DefaultRegistry().Kubectl(ctx, version) - if err != nil { - return nil, fmt.Errorf("no kubectl %s (%v)", version, err) - } - if installed { - p.logger.Info(fmt.Sprintf("kubectl %s has just been installed because of no pre-installed binary for that version", version)) - } - return NewKubectl(version, path), nil -} - -func (p *provider) findKustomize(ctx context.Context, version string) (*Kustomize, error) { - path, installed, err := toolregistry.DefaultRegistry().Kustomize(ctx, version) - if err != nil { - return nil, fmt.Errorf("no kustomize %s (%v)", version, err) - } - if installed { - p.logger.Info(fmt.Sprintf("kustomize %s has just been installed because of no pre-installed binary for that version", version)) - } - return NewKustomize(version, path, p.logger), nil -} - -func (p *provider) findHelm(ctx context.Context, version string) (*Helm, error) { - path, installed, err := toolregistry.DefaultRegistry().Helm(ctx, version) - if err != nil { - return nil, fmt.Errorf("no helm %s (%v)", version, err) - } - if installed { - p.logger.Info(fmt.Sprintf("helm %s has just been installed because of no pre-installed binary for that version", version)) - } - return NewHelm(version, path, p.logger), nil -} - -func determineTemplatingMethod(input config.KubernetesDeploymentInput, appDirPath string) TemplatingMethod { - if input.HelmChart != nil { - return TemplatingMethodHelm - } - if _, err := os.Stat(filepath.Join(appDirPath, kustomizationFileName)); err == nil { - return TemplatingMethodKustomize - } - return TemplatingMethodNone -} diff --git a/pkg/app/piped/cloudprovider/kubernetes/loader.go b/pkg/app/piped/cloudprovider/kubernetes/loader.go new file mode 100644 index 0000000000..4c8b00be0e --- /dev/null +++ b/pkg/app/piped/cloudprovider/kubernetes/loader.go @@ -0,0 +1,194 @@ +// Copyright 2022 The PipeCD 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 kubernetes + +import ( + "context" + "fmt" + "os" + "path/filepath" + "sync" + + "go.uber.org/zap" + + "github.com/pipe-cd/pipecd/pkg/app/piped/toolregistry" + "github.com/pipe-cd/pipecd/pkg/config" + "github.com/pipe-cd/pipecd/pkg/git" +) + +type TemplatingMethod string + +const ( + TemplatingMethodHelm TemplatingMethod = "helm" + TemplatingMethodKustomize TemplatingMethod = "kustomize" + TemplatingMethodNone TemplatingMethod = "none" +) + +type Loader interface { + // LoadManifests renders and loads all manifests for application. + LoadManifests(ctx context.Context) ([]Manifest, error) +} + +type gitClient interface { + Clone(ctx context.Context, repoID, remote, branch, destination string) (git.Repo, error) +} + +type loader struct { + appName string + appDir string + repoDir string + configFileName string + input config.KubernetesDeploymentInput + gc gitClient + logger *zap.Logger + + templatingMethod TemplatingMethod + kustomize *Kustomize + helm *Helm + initOnce sync.Once + initErr error +} + +func NewLoader( + appName, appDir, repoDir, configFileName string, + input config.KubernetesDeploymentInput, + gc gitClient, + logger *zap.Logger, +) Loader { + + return &loader{ + appName: appName, + appDir: appDir, + repoDir: repoDir, + configFileName: configFileName, + input: input, + gc: gc, + logger: logger.Named("kubernetes-loader"), + } +} + +// LoadManifests renders and loads all manifests for application. +func (l *loader) LoadManifests(ctx context.Context) (manifests []Manifest, err error) { + l.initOnce.Do(func() { + l.templatingMethod = determineTemplatingMethod(l.input, l.appDir) + switch l.templatingMethod { + case TemplatingMethodHelm: + l.helm, l.initErr = l.findHelm(ctx, l.input.HelmVersion) + + case TemplatingMethodKustomize: + l.kustomize, l.initErr = l.findKustomize(ctx, l.input.KustomizeVersion) + } + }) + if l.initErr != nil { + return nil, l.initErr + } + + switch l.templatingMethod { + case TemplatingMethodHelm: + var data string + switch { + case l.input.HelmChart.GitRemote != "": + chart := helmRemoteGitChart{ + GitRemote: l.input.HelmChart.GitRemote, + Ref: l.input.HelmChart.Ref, + Path: l.input.HelmChart.Path, + } + data, err = l.helm.TemplateRemoteGitChart(ctx, + l.appName, + l.appDir, + l.input.Namespace, + chart, + l.gc, + l.input.HelmOptions) + + case l.input.HelmChart.Repository != "": + chart := helmRemoteChart{ + Repository: l.input.HelmChart.Repository, + Name: l.input.HelmChart.Name, + Version: l.input.HelmChart.Version, + Insecure: l.input.HelmChart.Insecure, + } + data, err = l.helm.TemplateRemoteChart(ctx, + l.appName, + l.appDir, + l.input.Namespace, + chart, + l.input.HelmOptions) + + default: + data, err = l.helm.TemplateLocalChart(ctx, + l.appName, + l.appDir, + l.input.Namespace, + l.input.HelmChart.Path, + l.input.HelmOptions) + } + + if err != nil { + err = fmt.Errorf("unable to run helm template: %w", err) + return + } + manifests, err = ParseManifests(data) + + case TemplatingMethodKustomize: + var data string + data, err = l.kustomize.Template(ctx, l.appName, l.appDir, l.input.KustomizeOptions) + if err != nil { + err = fmt.Errorf("unable to run kustomize template: %w", err) + return + } + manifests, err = ParseManifests(data) + + case TemplatingMethodNone: + manifests, err = LoadPlainYAMLManifests(l.appDir, l.input.Manifests, l.configFileName) + + default: + err = fmt.Errorf("unsupport templating method %v", l.templatingMethod) + } + + return +} + +func (l *loader) findKustomize(ctx context.Context, version string) (*Kustomize, error) { + path, installed, err := toolregistry.DefaultRegistry().Kustomize(ctx, version) + if err != nil { + return nil, fmt.Errorf("no kustomize %s (%v)", version, err) + } + if installed { + l.logger.Info(fmt.Sprintf("kustomize %s has just been installed because of no pre-installed binary for that version", version)) + } + return NewKustomize(version, path, l.logger), nil +} + +func (l *loader) findHelm(ctx context.Context, version string) (*Helm, error) { + path, installed, err := toolregistry.DefaultRegistry().Helm(ctx, version) + if err != nil { + return nil, fmt.Errorf("no helm %s (%v)", version, err) + } + if installed { + l.logger.Info(fmt.Sprintf("helm %s has just been installed because of no pre-installed binary for that version", version)) + } + return NewHelm(version, path, l.logger), nil +} + +func determineTemplatingMethod(input config.KubernetesDeploymentInput, appDirPath string) TemplatingMethod { + if input.HelmChart != nil { + return TemplatingMethodHelm + } + if _, err := os.Stat(filepath.Join(appDirPath, kustomizationFileName)); err == nil { + return TemplatingMethodKustomize + } + return TemplatingMethodNone +} diff --git a/pkg/app/piped/driftdetector/kubernetes/detector.go b/pkg/app/piped/driftdetector/kubernetes/detector.go index a77e8e04bb..c44056b50b 100644 --- a/pkg/app/piped/driftdetector/kubernetes/detector.go +++ b/pkg/app/piped/driftdetector/kubernetes/detector.go @@ -250,7 +250,7 @@ func (d *detector) loadHeadManifests(ctx context.Context, app *model.Application } } - loader := provider.NewManifestLoader(app.Name, appDir, repoDir, app.GitPath.ConfigFilename, cfg.KubernetesApplicationSpec.Input, d.gitClient, d.logger) + loader := provider.NewLoader(app.Name, appDir, repoDir, app.GitPath.ConfigFilename, cfg.KubernetesApplicationSpec.Input, d.gitClient, d.logger) manifests, err = loader.LoadManifests(ctx) if err != nil { err = fmt.Errorf("failed to load new manifests: %w", err) diff --git a/pkg/app/piped/executor/kubernetes/baseline.go b/pkg/app/piped/executor/kubernetes/baseline.go index edf75a2015..3815658d45 100644 --- a/pkg/app/piped/executor/kubernetes/baseline.go +++ b/pkg/app/piped/executor/kubernetes/baseline.go @@ -85,7 +85,7 @@ func (e *deployExecutor) ensureBaselineRollout(ctx context.Context) model.StageS // Start rolling out the resources for BASELINE variant. e.LogPersister.Info("Start rolling out BASELINE variant...") - if err := applyManifests(ctx, e.provider, baselineManifests, e.appCfg.Input.Namespace, e.LogPersister); err != nil { + if err := applyManifests(ctx, e.applier, baselineManifests, e.appCfg.Input.Namespace, e.LogPersister); err != nil { return model.StageStatus_STAGE_FAILURE } @@ -101,7 +101,7 @@ func (e *deployExecutor) ensureBaselineClean(ctx context.Context) model.StageSta } resources := strings.Split(value, ",") - if err := removeBaselineResources(ctx, e.provider, resources, e.LogPersister); err != nil { + if err := removeBaselineResources(ctx, e.applier, resources, e.LogPersister); err != nil { e.LogPersister.Errorf("Unable to remove baseline resources: %v", err) return model.StageStatus_STAGE_FAILURE } diff --git a/pkg/app/piped/executor/kubernetes/canary.go b/pkg/app/piped/executor/kubernetes/canary.go index 7475f7a184..5f79ab24f2 100644 --- a/pkg/app/piped/executor/kubernetes/canary.go +++ b/pkg/app/piped/executor/kubernetes/canary.go @@ -47,7 +47,7 @@ func (e *deployExecutor) ensureCanaryRollout(ctx context.Context) model.StageSta e.Deployment.ApplicationId, e.commit, e.AppManifestsCache, - e.provider, + e.loader, e.Logger, ) if err != nil { @@ -102,7 +102,7 @@ func (e *deployExecutor) ensureCanaryRollout(ctx context.Context) model.StageSta // Start rolling out the resources for CANARY variant. e.LogPersister.Info("Start rolling out CANARY variant...") - if err := applyManifests(ctx, e.provider, canaryManifests, e.appCfg.Input.Namespace, e.LogPersister); err != nil { + if err := applyManifests(ctx, e.applier, canaryManifests, e.appCfg.Input.Namespace, e.LogPersister); err != nil { return model.StageStatus_STAGE_FAILURE } @@ -118,7 +118,7 @@ func (e *deployExecutor) ensureCanaryClean(ctx context.Context) model.StageStatu } resources := strings.Split(value, ",") - if err := removeCanaryResources(ctx, e.provider, resources, e.LogPersister); err != nil { + if err := removeCanaryResources(ctx, e.applier, resources, e.LogPersister); err != nil { e.LogPersister.Errorf("Unable to remove canary resources: %v", err) return model.StageStatus_STAGE_FAILURE } diff --git a/pkg/app/piped/executor/kubernetes/canary_test.go b/pkg/app/piped/executor/kubernetes/canary_test.go index 9c45178804..4a11d4f2b6 100644 --- a/pkg/app/piped/executor/kubernetes/canary_test.go +++ b/pkg/app/piped/executor/kubernetes/canary_test.go @@ -91,7 +91,7 @@ func TestEnsureCanaryRollout(t *testing.T) { }(), Logger: zap.NewNop(), }, - provider: func() provider.Provider { + loader: func() provider.Loader { p := providertest.NewMockProvider(ctrl) p.EXPECT().LoadManifests(gomock.Any()).Return(nil, fmt.Errorf("error")) return p @@ -122,7 +122,7 @@ func TestEnsureCanaryRollout(t *testing.T) { }(), Logger: zap.NewNop(), }, - provider: func() provider.Provider { + loader: func() provider.Loader { p := providertest.NewMockProvider(ctrl) p.EXPECT().LoadManifests(gomock.Any()).Return([]provider.Manifest{}, nil) return p @@ -155,7 +155,7 @@ func TestEnsureCanaryRollout(t *testing.T) { PipedConfig: &config.PipedSpec{}, Logger: zap.NewNop(), }, - provider: func() provider.Provider { + loader: func() provider.Loader { p := providertest.NewMockProvider(ctrl) p.EXPECT().LoadManifests(gomock.Any()).Return([]provider.Manifest{ provider.MakeManifest(provider.ResourceKey{ @@ -173,6 +173,10 @@ func TestEnsureCanaryRollout(t *testing.T) { }, }), }, nil) + return p + }(), + applier: func() provider.Applier { + p := providertest.NewMockProvider(ctrl) p.EXPECT().ApplyManifest(gomock.Any(), gomock.Any()).Return(fmt.Errorf("error")) return p }(), @@ -204,7 +208,7 @@ func TestEnsureCanaryRollout(t *testing.T) { PipedConfig: &config.PipedSpec{}, Logger: zap.NewNop(), }, - provider: func() provider.Provider { + loader: func() provider.Loader { p := providertest.NewMockProvider(ctrl) p.EXPECT().LoadManifests(gomock.Any()).Return([]provider.Manifest{ provider.MakeManifest(provider.ResourceKey{ @@ -222,6 +226,10 @@ func TestEnsureCanaryRollout(t *testing.T) { }, }), }, nil) + return p + }(), + applier: func() provider.Applier { + p := providertest.NewMockProvider(ctrl) p.EXPECT().ApplyManifest(gomock.Any(), gomock.Any()).Return(nil) return p }(), diff --git a/pkg/app/piped/executor/kubernetes/kubernetes.go b/pkg/app/piped/executor/kubernetes/kubernetes.go index a82edd9596..821a95218f 100644 --- a/pkg/app/piped/executor/kubernetes/kubernetes.go +++ b/pkg/app/piped/executor/kubernetes/kubernetes.go @@ -36,9 +36,11 @@ import ( type deployExecutor struct { executor.Input - commit string - appCfg *config.KubernetesApplicationSpec - provider provider.Provider + commit string + appCfg *config.KubernetesApplicationSpec + + loader provider.Loader + applier provider.Applier } type registerer interface { @@ -92,7 +94,19 @@ func (e *deployExecutor) Execute(sig executor.StopSignal) model.StageStatus { } } - e.provider = provider.NewProvider(e.Deployment.ApplicationName, ds.AppDir, ds.RepoDir, e.Deployment.GitPath.ConfigFilename, e.appCfg.Input, e.GitClient, e.Logger) + e.loader = provider.NewLoader( + e.Deployment.ApplicationName, + ds.AppDir, + ds.RepoDir, + e.Deployment.GitPath.ConfigFilename, + e.appCfg.Input, + e.GitClient, + e.Logger, + ) + e.applier = provider.NewApplier( + e.appCfg.Input, + e.Logger, + ) e.Logger.Info("start executing kubernetes stage", zap.String("stage-name", e.Stage.Name), zap.String("app-dir", ds.AppDir), @@ -147,7 +161,7 @@ func (e *deployExecutor) loadRunningManifests(ctx context.Context) (manifests [] return nil, err } - loader := provider.NewManifestLoader( + loader := provider.NewLoader( e.Deployment.ApplicationName, ds.AppDir, ds.RepoDir, @@ -171,7 +185,7 @@ func (l *manifestsLoadFunc) LoadManifests(ctx context.Context) ([]provider.Manif return l.loadFunc(ctx) } -func loadManifests(ctx context.Context, appID, commit string, manifestsCache cache.Cache, loader provider.ManifestLoader, logger *zap.Logger) (manifests []provider.Manifest, err error) { +func loadManifests(ctx context.Context, appID, commit string, manifestsCache cache.Cache, loader provider.Loader, logger *zap.Logger) (manifests []provider.Manifest, err error) { cache := provider.AppManifestsCache{ AppID: appID, Cache: manifestsCache, diff --git a/pkg/app/piped/executor/kubernetes/kubernetes_test.go b/pkg/app/piped/executor/kubernetes/kubernetes_test.go index 344b62b64c..365c0505ba 100644 --- a/pkg/app/piped/executor/kubernetes/kubernetes_test.go +++ b/pkg/app/piped/executor/kubernetes/kubernetes_test.go @@ -483,7 +483,7 @@ func TestDeleteResources(t *testing.T) { testcases := []struct { name string - provider provider.Provider + applier provider.Applier resources []provider.ResourceKey wantErr bool }{ @@ -500,7 +500,7 @@ func TestDeleteResources(t *testing.T) { Name: "foo", }, }, - provider: func() provider.Provider { + applier: func() provider.Applier { p := providertest.NewMockProvider(ctrl) p.EXPECT().Delete(gomock.Any(), gomock.Any()).Return(provider.ErrNotFound) return p @@ -514,7 +514,7 @@ func TestDeleteResources(t *testing.T) { Name: "foo", }, }, - provider: func() provider.Provider { + applier: func() provider.Applier { p := providertest.NewMockProvider(ctrl) p.EXPECT().Delete(gomock.Any(), gomock.Any()).Return(fmt.Errorf("unexpected error")) return p @@ -528,7 +528,7 @@ func TestDeleteResources(t *testing.T) { Name: "foo", }, }, - provider: func() provider.Provider { + applier: func() provider.Applier { p := providertest.NewMockProvider(ctrl) p.EXPECT().Delete(gomock.Any(), gomock.Any()).Return(nil) return p @@ -538,7 +538,7 @@ func TestDeleteResources(t *testing.T) { for _, tc := range testcases { t.Run(tc.name, func(t *testing.T) { ctx := context.Background() - err := deleteResources(ctx, tc.provider, tc.resources, &fakeLogPersister{}) + err := deleteResources(ctx, tc.applier, tc.resources, &fakeLogPersister{}) assert.Equal(t, tc.wantErr, err != nil) }) } diff --git a/pkg/app/piped/executor/kubernetes/primary.go b/pkg/app/piped/executor/kubernetes/primary.go index 6df314f3ac..16fc12fffa 100644 --- a/pkg/app/piped/executor/kubernetes/primary.go +++ b/pkg/app/piped/executor/kubernetes/primary.go @@ -42,7 +42,7 @@ func (e *deployExecutor) ensurePrimaryRollout(ctx context.Context) model.StageSt e.Deployment.ApplicationId, e.commit, e.AppManifestsCache, - e.provider, + e.loader, e.Logger, ) if err != nil { @@ -133,7 +133,7 @@ func (e *deployExecutor) ensurePrimaryRollout(ctx context.Context) model.StageSt // Start applying all manifests to add or update running resources. e.LogPersister.Info("Start rolling out PRIMARY variant...") - if err := applyManifests(ctx, e.provider, primaryManifests, e.appCfg.Input.Namespace, e.LogPersister); err != nil { + if err := applyManifests(ctx, e.applier, primaryManifests, e.appCfg.Input.Namespace, e.LogPersister); err != nil { return model.StageStatus_STAGE_FAILURE } e.LogPersister.Success("Successfully rolled out PRIMARY variant") @@ -171,7 +171,7 @@ func (e *deployExecutor) ensurePrimaryRollout(ctx context.Context) model.StageSt // Start deleting all running resources that are not defined in Git. e.LogPersister.Infof("Start deleting %d resources", len(removeKeys)) - if err := deleteResources(ctx, e.provider, removeKeys, e.LogPersister); err != nil { + if err := deleteResources(ctx, e.applier, removeKeys, e.LogPersister); err != nil { return model.StageStatus_STAGE_FAILURE } diff --git a/pkg/app/piped/executor/kubernetes/primary_test.go b/pkg/app/piped/executor/kubernetes/primary_test.go index d2955b170d..7b0917b907 100644 --- a/pkg/app/piped/executor/kubernetes/primary_test.go +++ b/pkg/app/piped/executor/kubernetes/primary_test.go @@ -91,7 +91,7 @@ func TestEnsurePrimaryRollout(t *testing.T) { }(), Logger: zap.NewNop(), }, - provider: func() provider.Provider { + loader: func() provider.Loader { p := providertest.NewMockProvider(ctrl) p.EXPECT().LoadManifests(gomock.Any()).Return(nil, fmt.Errorf("error")) return p @@ -125,7 +125,7 @@ func TestEnsurePrimaryRollout(t *testing.T) { }(), Logger: zap.NewNop(), }, - provider: func() provider.Provider { + loader: func() provider.Loader { p := providertest.NewMockProvider(ctrl) p.EXPECT().LoadManifests(gomock.Any()).Return([]provider.Manifest{ provider.MakeManifest(provider.ResourceKey{ @@ -135,6 +135,10 @@ func TestEnsurePrimaryRollout(t *testing.T) { Object: map[string]interface{}{"spec": map[string]interface{}{}}, }), }, nil) + return p + }(), + applier: func() provider.Applier { + p := providertest.NewMockProvider(ctrl) p.EXPECT().ApplyManifest(gomock.Any(), gomock.Any()).Return(nil) return p }(), @@ -167,7 +171,7 @@ func TestEnsurePrimaryRollout(t *testing.T) { }(), Logger: zap.NewNop(), }, - provider: func() provider.Provider { + loader: func() provider.Loader { p := providertest.NewMockProvider(ctrl) p.EXPECT().LoadManifests(gomock.Any()).Return([]provider.Manifest{ provider.MakeManifest(provider.ResourceKey{ @@ -184,6 +188,10 @@ func TestEnsurePrimaryRollout(t *testing.T) { Object: map[string]interface{}{"spec": map[string]interface{}{}}, }), }, nil) + return p + }(), + applier: func() provider.Applier { + p := providertest.NewMockProvider(ctrl) p.EXPECT().ApplyManifest(gomock.Any(), gomock.Any()).Return(nil) p.EXPECT().ApplyManifest(gomock.Any(), gomock.Any()).Return(nil) return p @@ -222,7 +230,7 @@ func TestEnsurePrimaryRollout(t *testing.T) { }(), Logger: zap.NewNop(), }, - provider: func() provider.Provider { + loader: func() provider.Loader { p := providertest.NewMockProvider(ctrl) p.EXPECT().LoadManifests(gomock.Any()).Return([]provider.Manifest{ provider.MakeManifest(provider.ResourceKey{ @@ -267,7 +275,7 @@ func TestEnsurePrimaryRollout(t *testing.T) { }(), Logger: zap.NewNop(), }, - provider: func() provider.Provider { + loader: func() provider.Loader { p := providertest.NewMockProvider(ctrl) p.EXPECT().LoadManifests(gomock.Any()).Return([]provider.Manifest{ provider.MakeManifest(provider.ResourceKey{ diff --git a/pkg/app/piped/executor/kubernetes/rollback.go b/pkg/app/piped/executor/kubernetes/rollback.go index 7c8d374617..1e0fa451a3 100644 --- a/pkg/app/piped/executor/kubernetes/rollback.go +++ b/pkg/app/piped/executor/kubernetes/rollback.go @@ -74,7 +74,7 @@ func (e *rollbackExecutor) ensureRollback(ctx context.Context) model.StageStatus } } - p := provider.NewProvider(e.Deployment.ApplicationName, ds.AppDir, ds.RepoDir, e.Deployment.GitPath.ConfigFilename, appCfg.Input, e.GitClient, e.Logger) + loader := provider.NewLoader(e.Deployment.ApplicationName, ds.AppDir, ds.RepoDir, e.Deployment.GitPath.ConfigFilename, appCfg.Input, e.GitClient, e.Logger) e.Logger.Info("start executing kubernetes stage", zap.String("stage-name", e.Stage.Name), zap.String("app-dir", ds.AppDir), @@ -85,7 +85,14 @@ func (e *rollbackExecutor) ensureRollback(ctx context.Context) model.StageStatus // Load the manifests at the specified commit. e.LogPersister.Infof("Loading manifests at running commit %s for handling", e.Deployment.RunningCommitHash) - manifests, err := loadManifests(ctx, e.Deployment.ApplicationId, e.Deployment.RunningCommitHash, e.AppManifestsCache, p, e.Logger) + manifests, err := loadManifests( + ctx, + e.Deployment.ApplicationId, + e.Deployment.RunningCommitHash, + e.AppManifestsCache, + loader, + e.Logger, + ) if err != nil { e.LogPersister.Errorf("Failed while loading running manifests (%v)", err) return model.StageStatus_STAGE_FAILURE @@ -128,8 +135,13 @@ func (e *rollbackExecutor) ensureRollback(ctx context.Context) model.StageStatus return model.StageStatus_STAGE_FAILURE } + applier := provider.NewApplier( + appCfg.Input, + e.Logger, + ) + // Start applying all manifests to add or update running resources. - if err := applyManifests(ctx, p, manifests, appCfg.Input.Namespace, e.LogPersister); err != nil { + if err := applyManifests(ctx, applier, manifests, appCfg.Input.Namespace, e.LogPersister); err != nil { return model.StageStatus_STAGE_FAILURE } @@ -139,7 +151,7 @@ func (e *rollbackExecutor) ensureRollback(ctx context.Context) model.StageStatus e.LogPersister.Info("Start checking to ensure that the CANARY variant should be removed") if value, ok := e.MetadataStore.Shared().Get(addedCanaryResourcesMetadataKey); ok { resources := strings.Split(value, ",") - if err := removeCanaryResources(ctx, p, resources, e.LogPersister); err != nil { + if err := removeCanaryResources(ctx, applier, resources, e.LogPersister); err != nil { errs = append(errs, err) } } @@ -148,7 +160,7 @@ func (e *rollbackExecutor) ensureRollback(ctx context.Context) model.StageStatus e.LogPersister.Info("Start checking to ensure that the BASELINE variant should be removed") if value, ok := e.MetadataStore.Shared().Get(addedBaselineResourcesMetadataKey); ok { resources := strings.Split(value, ",") - if err := removeBaselineResources(ctx, p, resources, e.LogPersister); err != nil { + if err := removeBaselineResources(ctx, applier, resources, e.LogPersister); err != nil { errs = append(errs, err) } } diff --git a/pkg/app/piped/executor/kubernetes/sync.go b/pkg/app/piped/executor/kubernetes/sync.go index c8632c4b7b..86c734ada4 100644 --- a/pkg/app/piped/executor/kubernetes/sync.go +++ b/pkg/app/piped/executor/kubernetes/sync.go @@ -30,7 +30,7 @@ func (e *deployExecutor) ensureSync(ctx context.Context) model.StageStatus { e.Deployment.ApplicationId, e.commit, e.AppManifestsCache, - e.provider, + e.loader, e.Logger, ) if err != nil { @@ -76,7 +76,7 @@ func (e *deployExecutor) ensureSync(ctx context.Context) model.StageStatus { } // Start applying all manifests to add or update running resources. - if err := applyManifests(ctx, e.provider, manifests, e.appCfg.Input.Namespace, e.LogPersister); err != nil { + if err := applyManifests(ctx, e.applier, manifests, e.appCfg.Input.Namespace, e.LogPersister); err != nil { return model.StageStatus_STAGE_FAILURE } @@ -112,7 +112,7 @@ func (e *deployExecutor) ensureSync(ctx context.Context) model.StageStatus { e.LogPersister.Infof("Found %d live resources that are no longer defined in Git", len(removeKeys)) // Start deleting all running resources that are not defined in Git. - if err := deleteResources(ctx, e.provider, removeKeys, e.LogPersister); err != nil { + if err := deleteResources(ctx, e.applier, removeKeys, e.LogPersister); err != nil { return model.StageStatus_STAGE_FAILURE } diff --git a/pkg/app/piped/executor/kubernetes/sync_test.go b/pkg/app/piped/executor/kubernetes/sync_test.go index 465f81d5df..2c637ac3bd 100644 --- a/pkg/app/piped/executor/kubernetes/sync_test.go +++ b/pkg/app/piped/executor/kubernetes/sync_test.go @@ -62,7 +62,7 @@ func TestEnsureSync(t *testing.T) { }(), Logger: zap.NewNop(), }, - provider: func() provider.Provider { + loader: func() provider.Loader { p := providertest.NewMockProvider(ctrl) p.EXPECT().LoadManifests(gomock.Any()).Return(nil, fmt.Errorf("error")) return p @@ -89,7 +89,7 @@ func TestEnsureSync(t *testing.T) { }(), Logger: zap.NewNop(), }, - provider: func() provider.Provider { + loader: func() provider.Loader { p := providertest.NewMockProvider(ctrl) p.EXPECT().LoadManifests(gomock.Any()).Return([]provider.Manifest{ provider.MakeManifest(provider.ResourceKey{ @@ -99,6 +99,10 @@ func TestEnsureSync(t *testing.T) { Object: map[string]interface{}{"spec": map[string]interface{}{}}, }), }, nil) + return p + }(), + applier: func() provider.Applier { + p := providertest.NewMockProvider(ctrl) p.EXPECT().ApplyManifest(gomock.Any(), gomock.Any()).Return(fmt.Errorf("error")) return p }(), @@ -129,7 +133,7 @@ func TestEnsureSync(t *testing.T) { }(), Logger: zap.NewNop(), }, - provider: func() provider.Provider { + loader: func() provider.Loader { p := providertest.NewMockProvider(ctrl) p.EXPECT().LoadManifests(gomock.Any()).Return([]provider.Manifest{ provider.MakeManifest(provider.ResourceKey{ @@ -139,6 +143,10 @@ func TestEnsureSync(t *testing.T) { Object: map[string]interface{}{"spec": map[string]interface{}{}}, }), }, nil) + return p + }(), + applier: func() provider.Applier { + p := providertest.NewMockProvider(ctrl) p.EXPECT().ApplyManifest(gomock.Any(), gomock.Any()).Return(nil) return p }(), diff --git a/pkg/app/piped/executor/kubernetes/traffic.go b/pkg/app/piped/executor/kubernetes/traffic.go index 8e599d3e65..89535eba51 100644 --- a/pkg/app/piped/executor/kubernetes/traffic.go +++ b/pkg/app/piped/executor/kubernetes/traffic.go @@ -56,7 +56,7 @@ func (e *deployExecutor) ensureTrafficRouting(ctx context.Context) model.StageSt e.Deployment.ApplicationId, e.commit, e.AppManifestsCache, - e.provider, + e.loader, e.Logger, ) if err != nil { @@ -134,7 +134,7 @@ func (e *deployExecutor) ensureTrafficRouting(ctx context.Context) model.StageSt canaryPercent, baselinePercent, ) - if err := applyManifests(ctx, e.provider, []provider.Manifest{trafficRoutingManifest}, e.appCfg.Input.Namespace, e.LogPersister); err != nil { + if err := applyManifests(ctx, e.applier, []provider.Manifest{trafficRoutingManifest}, e.appCfg.Input.Namespace, e.LogPersister); err != nil { return model.StageStatus_STAGE_FAILURE } diff --git a/pkg/app/piped/planner/kubernetes/kubernetes.go b/pkg/app/piped/planner/kubernetes/kubernetes.go index 0ccca05d76..b9a10fa9d3 100644 --- a/pkg/app/piped/planner/kubernetes/kubernetes.go +++ b/pkg/app/piped/planner/kubernetes/kubernetes.go @@ -81,7 +81,7 @@ func (p *Planner) Plan(ctx context.Context, in planner.Input) (out planner.Outpu newManifests, ok := manifestCache.Get(in.Trigger.Commit.Hash) if !ok { // When the manifests were not in the cache we have to load them. - loader := provider.NewManifestLoader(in.ApplicationName, ds.AppDir, ds.RepoDir, in.GitPath.ConfigFilename, cfg.Input, in.GitClient, in.Logger) + loader := provider.NewLoader(in.ApplicationName, ds.AppDir, ds.RepoDir, in.GitPath.ConfigFilename, cfg.Input, in.GitClient, in.Logger) newManifests, err = loader.LoadManifests(ctx) if err != nil { return @@ -200,7 +200,7 @@ func (p *Planner) Plan(ctx context.Context, in planner.Input) (out planner.Outpu return } - loader := provider.NewManifestLoader(in.ApplicationName, runningDs.AppDir, runningDs.RepoDir, in.GitPath.ConfigFilename, cfg.Input, in.GitClient, in.Logger) + loader := provider.NewLoader(in.ApplicationName, runningDs.AppDir, runningDs.RepoDir, in.GitPath.ConfigFilename, cfg.Input, in.GitClient, in.Logger) oldManifests, err = loader.LoadManifests(ctx) if err != nil { err = fmt.Errorf("failed to load previously deployed manifests: %w", err) diff --git a/pkg/app/piped/planpreview/kubernetesdiff.go b/pkg/app/piped/planpreview/kubernetesdiff.go index f2d827789b..38d5ac1288 100644 --- a/pkg/app/piped/planpreview/kubernetesdiff.go +++ b/pkg/app/piped/planpreview/kubernetesdiff.go @@ -116,7 +116,7 @@ func loadKubernetesManifests(ctx context.Context, app model.Application, dsp dep return nil, fmt.Errorf("malformed application configuration file") } - loader := provider.NewManifestLoader( + loader := provider.NewLoader( app.Name, ds.AppDir, ds.RepoDir,