diff --git a/pkg/apis/core/v1/types.go b/pkg/apis/core/v1/types.go index fe9db625..408402fe 100644 --- a/pkg/apis/core/v1/types.go +++ b/pkg/apis/core/v1/types.go @@ -19,6 +19,7 @@ package v1 import ( apiv1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime/schema" ) // @@ -675,4 +676,11 @@ type ( CanaryConfigStatus struct { Status string `json:"status"` } + + // MetadataAccessor lets you work with object metadata and type metadata + // from any of the versioned or internal API objects. + MetadataAccessor interface { + GetObjectKind() schema.ObjectKind + GetObjectMeta() metav1.Object + } ) diff --git a/pkg/fission-cli/cmd/spec/apply.go b/pkg/fission-cli/cmd/spec/apply.go index ab9a3b01..f69409ff 100644 --- a/pkg/fission-cli/cmd/spec/apply.go +++ b/pkg/fission-cli/cmd/spec/apply.go @@ -115,14 +115,9 @@ func (opts *ApplySubCommand) run(input cli.Input) error { return errors.Wrap(err, "error reading specs") } - var warnings []string - // validate - warnings, err = fr.Validate(input) + err = Validate(input) if err != nil { - return errors.Wrap(err, "error validating specs") - } - for _, warning := range warnings { - console.Warn(warning) + return errors.Wrap(err, "abort applying resources") } // make changes to the cluster based on the specs diff --git a/pkg/fission-cli/cmd/spec/list.go b/pkg/fission-cli/cmd/spec/list.go index 0c165298..7e7c7c11 100644 --- a/pkg/fission-cli/cmd/spec/list.go +++ b/pkg/fission-cli/cmd/spec/list.go @@ -26,6 +26,7 @@ import ( "github.com/pkg/errors" fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/controller/client" "github.com/fission/fission/pkg/fission-cli/cliwrapper/cli" "github.com/fission/fission/pkg/fission-cli/cmd" flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" @@ -58,56 +59,56 @@ func (opts *ListSubCommand) run(input cli.Input) error { deployID = fr.DeploymentConfig.UID } - allfn, err := getAllFunctions(opts) + allfn, err := getAllFunctions(opts.Client()) if err != nil { return errors.Wrap(err, "error getting Functions from all namespaces") } specfns := getAppliedFunctions(allfn, deployID) ShowFunctions(specfns) - allenvs, err := getAllEnvironments(opts) + allenvs, err := getAllEnvironments(opts.Client()) if err != nil { return errors.Wrap(err, "error getting Environments from all namespaces") } specenvs := getAppliedEnvironments(allenvs, deployID) ShowEnvironments(specenvs) - pkglists, err := getAllPackages(opts) + pkglists, err := getAllPackages(opts.Client()) if err != nil { return errors.Wrap(err, "error getting Packages from all namespaces") } specPkgs := getAppliedPackages(pkglists, deployID) ShowPackages(specPkgs) - canaryCfgs, err := getAllCanaryConfigs(opts) + canaryCfgs, err := getAllCanaryConfigs(opts.Client()) if err != nil { return errors.Wrap(err, "error getting Canary Config from all namespaces") } specCanaryCfgs := getAppliedCanaryConfigs(canaryCfgs, deployID) ShowCanaryConfigs(specCanaryCfgs) - hts, err := getAllHTTPTriggers(opts) + hts, err := getAllHTTPTriggers(opts.Client()) if err != nil { return errors.Wrap(err, "error getting HTTP Triggers from all namespaces") } specHTTPTriggers := getAppliedHTTPTriggers(hts, deployID) ShowHTTPTriggers(specHTTPTriggers) - mqts, err := getAllMessageQueueTriggers(opts, input.String(flagkey.MqtMQType)) + mqts, err := getAllMessageQueueTriggers(opts.Client(), input.String(flagkey.MqtMQType)) if err != nil { return errors.Wrap(err, "error getting MessageQueue Triggers from all namespaces") } specMessageQueueTriggers := getAppliedMessageQueueTriggers(mqts, deployID) ShowMQTriggers(specMessageQueueTriggers) - tts, err := getAllTimeTriggers(opts) + tts, err := getAllTimeTriggers(opts.Client()) if err != nil { return errors.Wrap(err, "error getting Time Triggers from all namespaces") } specTimeTriggers := getAppliedTimeTriggers(tts, deployID) ShowTimeTriggers(specTimeTriggers) - kws, err := getAllKubeWatchTriggers(opts) + kws, err := getAllKubeWatchTriggers(opts.Client()) if err != nil { return errors.Wrap(err, "error getting Kube Watchers from all namespaces") } @@ -382,8 +383,8 @@ func ShowAppliedKubeWatchers(ws []fv1.KubernetesWatchTrigger) { } // getAllFunctions get lists of functions in all namespaces -func getAllFunctions(opts *ListSubCommand) ([]fv1.Function, error) { - fns, err := opts.Client().V1().Function().List("") +func getAllFunctions(client client.Interface) ([]fv1.Function, error) { + fns, err := client.V1().Function().List("") if err != nil { return nil, errors.Errorf("Unable to get Functions %v", err.Error()) } @@ -391,17 +392,17 @@ func getAllFunctions(opts *ListSubCommand) ([]fv1.Function, error) { } // getAllEnvironments get lists of environments in all namespaces -func getAllEnvironments(opts *ListSubCommand) ([]fv1.Environment, error) { - envs, err := opts.Client().V1().Environment().List("") +func getAllEnvironments(client client.Interface) ([]fv1.Environment, error) { + envs, err := client.V1().Environment().List("") if err != nil { - return nil, errors.Errorf("Unable to get Enviornmets %v", err.Error()) + return nil, errors.Errorf("Unable to get Environments %v", err.Error()) } return envs, nil } // getAllPackages get lists of packages in all namespaces -func getAllPackages(opts *ListSubCommand) ([]fv1.Package, error) { - pkgList, err := opts.Client().V1().Package().List("") +func getAllPackages(client client.Interface) ([]fv1.Package, error) { + pkgList, err := client.V1().Package().List("") if err != nil { return nil, errors.Errorf("Unable to get Packages %v", err.Error()) } @@ -409,44 +410,44 @@ func getAllPackages(opts *ListSubCommand) ([]fv1.Package, error) { } // getAllCanaryConfigs get lists of canary configs in all namespaces -func getAllCanaryConfigs(opts *ListSubCommand) ([]fv1.CanaryConfig, error) { - canaryCfgs, err := opts.Client().V1().CanaryConfig().List("") +func getAllCanaryConfigs(client client.Interface) ([]fv1.CanaryConfig, error) { + canaryCfgs, err := client.V1().CanaryConfig().List("") if err != nil { return nil, errors.Errorf("Unable to get Canary Configs %v", err.Error()) } return canaryCfgs, nil } -// getAllHTTPTriggers get lists of HTTP Triggers in all namespaces -func getAllHTTPTriggers(opts *ListSubCommand) ([]fv1.HTTPTrigger, error) { - hts, err := opts.Client().V1().HTTPTrigger().List("") +// getAllHTTPTriggers get lists of HTTP Triggers in all namespaces +func getAllHTTPTriggers(client client.Interface) ([]fv1.HTTPTrigger, error) { + hts, err := client.V1().HTTPTrigger().List("") if err != nil { return nil, errors.Errorf("Unable to get HTTP Triggers %v", err.Error()) } return hts, nil } -// getAllMessageQueueTriggers get lists of MessageQueue Triggers in all namespaces -func getAllMessageQueueTriggers(opts *ListSubCommand, mqttype string) ([]fv1.MessageQueueTrigger, error) { - mqts, err := opts.Client().V1().MessageQueueTrigger().List(mqttype, "") +// getAllMessageQueueTriggers get lists of MessageQueue Triggers in all namespaces +func getAllMessageQueueTriggers(client client.Interface, mqttype string) ([]fv1.MessageQueueTrigger, error) { + mqts, err := client.V1().MessageQueueTrigger().List(mqttype, "") if err != nil { return nil, errors.Errorf("Unable to get MessageQueue Triggers %v", err.Error()) } return mqts, nil } -// getAllTimeTriggers get lists of Time Triggers in all namespaces -func getAllTimeTriggers(opts *ListSubCommand) ([]fv1.TimeTrigger, error) { - tts, err := opts.Client().V1().TimeTrigger().List("") +// getAllTimeTriggers get lists of Time Triggers in all namespaces +func getAllTimeTriggers(client client.Interface) ([]fv1.TimeTrigger, error) { + tts, err := client.V1().TimeTrigger().List("") if err != nil { return nil, errors.Errorf("Unable to get Time Triggers %v", err.Error()) } return tts, nil } -// getAllKubeWatchTriggers get lists of Kube Watchers in all namespaces -func getAllKubeWatchTriggers(opts *ListSubCommand) ([]fv1.KubernetesWatchTrigger, error) { - ws, err := opts.Client().V1().KubeWatcher().List("") +// getAllKubeWatchTriggers get lists of Kube Watchers in all namespaces +func getAllKubeWatchTriggers(client client.Interface) ([]fv1.KubernetesWatchTrigger, error) { + ws, err := client.V1().KubeWatcher().List("") if err != nil { return nil, errors.Errorf("Unable to get Kube Watchers %v", err.Error()) } diff --git a/pkg/fission-cli/cmd/spec/validate.go b/pkg/fission-cli/cmd/spec/validate.go index 82761702..96a431e5 100644 --- a/pkg/fission-cli/cmd/spec/validate.go +++ b/pkg/fission-cli/cmd/spec/validate.go @@ -28,10 +28,12 @@ import ( "github.com/pkg/errors" fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/controller/client" "github.com/fission/fission/pkg/fission-cli/cliwrapper/cli" "github.com/fission/fission/pkg/fission-cli/cmd" "github.com/fission/fission/pkg/fission-cli/console" "github.com/fission/fission/pkg/fission-cli/util" + "github.com/fission/fission/pkg/utils" ) type ValidateSubCommand struct { @@ -57,18 +59,140 @@ func (opts *ValidateSubCommand) run(input cli.Input) error { return errors.Wrap(err, "error reading specs") } + console.Infof("DeployUID: %v", fr.DeploymentConfig.UID) + console.Infof("Resources:\n * %v Functions\n * %v Environments\n * %v Packages \n * %v Http Triggers \n * %v MessageQueue Triggers\n * %v Time Triggers\n * %v Kube Watchers\n * %v ArchiveUploadSpec\n", + len(fr.Functions), len(fr.Environments), len(fr.Packages), len(fr.HttpTriggers), len(fr.MessageQueueTriggers), len(fr.TimeTriggers), len(fr.KubernetesWatchTriggers), len(fr.ArchiveUploadSpecs)) + var warnings []string // this does the rest of the checks, like dangling refs warnings, err = fr.Validate(input) if err != nil { return errors.Wrap(err, "error validating specs") } - fmt.Printf("Spec validation successful\nSpec contains\n %v Functions\n %v Environments\n %v Packages \n %v Http Triggers \n %v MessageQueue Triggers\n %v Time Triggers\n %v Kube Watchers\n %v ArchiveUploadSpec\n", - len(fr.Functions), len(fr.Environments), len(fr.Packages), len(fr.HttpTriggers), len(fr.MessageQueueTriggers), len(fr.TimeTriggers), len(fr.KubernetesWatchTriggers), len(fr.ArchiveUploadSpecs)) + + err = resourceConflictCheck(opts.Client(), fr) + if err != nil { + return errors.Wrap(err, "name conflict error") + } for _, warning := range warnings { console.Warn(warning) } + + console.Info("Validation Successful") + + return nil +} + +// resourceConflictCheck checks if any of the spec resources with +// the same name is already present in the same cluster namespace. +// If a same name resource exists in the same namespace, a name +// conflict error will be returned. +func resourceConflictCheck(c client.Interface, fr *FissionResources) error { + deployUID := fr.DeploymentConfig.UID + result := utils.MultiErrorWithFormat() + + fnList, err := getAllFunctions(c) + if err != nil { + return errors.Errorf("Unable to get Functions %v", err.Error()) + } + for _, sObj := range fr.Functions { + for _, cObj := range fnList { + if err := isResourceConflicts(deployUID, &sObj, &cObj); err != nil { + result = multierror.Append(result, err) + break + } + } + } + + envList, err := getAllEnvironments(c) + if err != nil { + return errors.Errorf("Unable to get Environments %v", err.Error()) + } + for _, sObj := range fr.Environments { + for _, cObj := range envList { + if err := isResourceConflicts(deployUID, &sObj, &cObj); err != nil { + result = multierror.Append(result, err) + break + } + } + } + + pkgList, err := getAllPackages(c) + if err != nil { + return errors.Errorf("Unable to get Packages %v", err.Error()) + } + for _, sObj := range fr.Packages { + for _, cObj := range pkgList { + if err := isResourceConflicts(deployUID, &sObj, &cObj); err != nil { + result = multierror.Append(result, err) + break + } + } + } + + httptriggerList, err := getAllHTTPTriggers(c) + if err != nil { + return errors.Errorf("Unable to get HTTPTrigger %v", err.Error()) + } + for _, sObj := range fr.HttpTriggers { + for _, cObj := range httptriggerList { + if err := isResourceConflicts(deployUID, &sObj, &cObj); err != nil { + result = multierror.Append(result, err) + break + } + } + } + + mqtriggerList, err := getAllMessageQueueTriggers(c, "") + if err != nil { + return errors.Errorf("Unable to get Message Queue Trigger %v", err.Error()) + } + for _, sObj := range fr.MessageQueueTriggers { + for _, cObj := range mqtriggerList { + if err := isResourceConflicts(deployUID, &sObj, &cObj); err != nil { + result = multierror.Append(result, err) + break + } + } + } + + timetriggerList, err := getAllTimeTriggers(c) + if err != nil { + return errors.Errorf("Unable to get Time Trigger %v", err.Error()) + } + for _, sObj := range fr.TimeTriggers { + for _, cObj := range timetriggerList { + if err := isResourceConflicts(deployUID, &sObj, &cObj); err != nil { + result = multierror.Append(result, err) + break + } + } + } + + kubewatchtriggerList, err := getAllKubeWatchTriggers(c) + if err != nil { + return errors.Errorf("Unable to get Kubernetes Watch Trigger %v", err.Error()) + } + for _, sObj := range fr.KubernetesWatchTriggers { + for _, cObj := range kubewatchtriggerList { + if err := isResourceConflicts(deployUID, &sObj, &cObj); err != nil { + result = multierror.Append(result, err) + break + } + } + } + + return result.ErrorOrNil() +} + +func isResourceConflicts(deployUID string, specObj fv1.MetadataAccessor, clusterObj fv1.MetadataAccessor) error { + if specObj.GetObjectMeta().GetName() == clusterObj.GetObjectMeta().GetName() && + specObj.GetObjectMeta().GetNamespace() == clusterObj.GetObjectMeta().GetNamespace() && + deployUID != clusterObj.GetObjectMeta().GetAnnotations()[FISSION_DEPLOYMENT_UID_KEY] { + return fmt.Errorf("%v: '%v/%v' with different deploy uid already exists", + clusterObj.GetObjectKind().GroupVersionKind().Kind, clusterObj.GetObjectMeta().GetName(), clusterObj.GetObjectMeta().GetNamespace()) + } return nil } diff --git a/pkg/fission-cli/console/log.go b/pkg/fission-cli/console/log.go index e4373113..15a6556d 100644 --- a/pkg/fission-cli/console/log.go +++ b/pkg/fission-cli/console/log.go @@ -43,7 +43,11 @@ func Warn(msg interface{}) { } func Info(msg interface{}) { - os.Stderr.WriteString(fmt.Sprintf("%v\n", trimNewline(msg))) + os.Stdout.WriteString(fmt.Sprintf("%v\n", trimNewline(msg))) +} + +func Infof(format string, args ...interface{}) { + os.Stdout.WriteString(fmt.Sprintf("%v\n", trimNewline(fmt.Sprintf(format, args...)))) } func Verbose(verbosityLevel int, format string, args ...interface{}) {