diff --git a/.golangci.yaml b/.golangci.yaml index 58f3efc2..1a28e90b 100644 --- a/.golangci.yaml +++ b/.golangci.yaml @@ -1,7 +1,22 @@ +linters: + enable: + - deadcode + - gofmt + - goimports + - gosimple + - govet + - ineffassign + - misspell + - nakedret + - staticcheck + - structcheck + - typecheck + - unconvert + - varcheck linters-settings: errcheck: ignore: go.uber.org/zap:Sync goimports: # put imports beginning with prefix after 3rd-party packages; # it's a comma-separated list of prefixes - local-prefixes: github.com/trussworks/my-cli-tool \ No newline at end of file + local-prefixes: github.com/fission/fission \ No newline at end of file diff --git a/cmd/reporter/app/cmd_event.go b/cmd/reporter/app/cmd_event.go index 74477a72..431f35c6 100644 --- a/cmd/reporter/app/cmd_event.go +++ b/cmd/reporter/app/cmd_event.go @@ -18,8 +18,9 @@ package app import ( "log" - "github.com/fission/fission/pkg/tracker" "github.com/spf13/cobra" + + "github.com/fission/fission/pkg/tracker" ) func eventCommandHandler(cmd *cobra.Command, args []string) error { diff --git a/pkg/builder/client/client.go b/pkg/builder/client/client.go index a70e9aaa..f23b0d26 100644 --- a/pkg/builder/client/client.go +++ b/pkg/builder/client/client.go @@ -82,7 +82,7 @@ func (c *Client) Build(req *builder.PackageBuildRequest) (*builder.PackageBuildR } pkgBuildResp := builder.PackageBuildResponse{} - err = json.Unmarshal([]byte(rBody), &pkgBuildResp) + err = json.Unmarshal(rBody, &pkgBuildResp) if err != nil { c.logger.Error("error parsing resp body", zap.Error(err)) return nil, err diff --git a/pkg/controller/packageApi.go b/pkg/controller/packageApi.go index d5f031bb..d2edaef7 100644 --- a/pkg/controller/packageApi.go +++ b/pkg/controller/packageApi.go @@ -189,7 +189,7 @@ func (a *API) PackageApiGet(w http.ResponseWriter, r *http.Request) { var resp []byte if raw != "" { - resp = []byte(f.Spec.Deployment.Literal) + resp = f.Spec.Deployment.Literal } else { resp, err = json.Marshal(f) if err != nil { diff --git a/pkg/executor/executortype/newdeploy/newdeploy.go b/pkg/executor/executortype/newdeploy/newdeploy.go index 32c94822..e2b1d0e1 100644 --- a/pkg/executor/executortype/newdeploy/newdeploy.go +++ b/pkg/executor/executortype/newdeploy/newdeploy.go @@ -48,7 +48,7 @@ const ( func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Environment, deployName string, deployLabels map[string]string, deployAnnotations map[string]string, deployNamespace string) (*appsv1.Deployment, error) { - specializationTimeout := int(fn.Spec.InvokeStrategy.ExecutionStrategy.SpecializationTimeout) + specializationTimeout := fn.Spec.InvokeStrategy.ExecutionStrategy.SpecializationTimeout minScale := int32(fn.Spec.InvokeStrategy.ExecutionStrategy.MinScale) // Always scale to at least one pod when createOrGetDeployment diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index a75c82d2..14d21ebc 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -207,7 +207,7 @@ func (deploy *NewDeploy) getServiceInfo(obj apiv1.ObjectReference) (*apiv1.Servi if err != nil || !exists { deploy.logger.Debug( - "Falling back to getting service info from k8s API -- this may cause performace issues for your function.", + "Falling back to getting service info from k8s API -- this may cause performance issues for your function.", zap.Bool("exists", exists), zap.Error(err), ) @@ -224,7 +224,7 @@ func (deploy *NewDeploy) getDeploymentInfo(obj apiv1.ObjectReference) (*appsv1.D if err != nil || !exists { deploy.logger.Debug( - "Falling back to getting deployment info from k8s API -- this may cause performace issues for your function.", + "Falling back to getting deployment info from k8s API -- this may cause performance issues for your function.", zap.Bool("exists", exists), zap.Error(err), ) diff --git a/pkg/executor/executortype/poolmgr/gp.go b/pkg/executor/executortype/poolmgr/gp.go index 890dfcc0..1c506ed2 100644 --- a/pkg/executor/executortype/poolmgr/gp.go +++ b/pkg/executor/executortype/poolmgr/gp.go @@ -28,7 +28,6 @@ import ( "time" "github.com/dchest/uniuri" - "github.com/fission/fission/pkg/utils" "github.com/pkg/errors" "go.uber.org/zap" appsv1 "k8s.io/api/apps/v1" @@ -50,6 +49,7 @@ import ( "github.com/fission/fission/pkg/executor/util" fetcherClient "github.com/fission/fission/pkg/fetcher/client" fetcherConfig "github.com/fission/fission/pkg/fetcher/config" + "github.com/fission/fission/pkg/utils" ) type ( @@ -171,32 +171,58 @@ func (gp *GenericPool) getDeployAnnotations() map[string]string { } } +func (gp *GenericPool) checkMetricsApi() bool { + apiGroups, err := gp.metricsClient.DiscoveryClient.ServerGroups() + if err != nil { + gp.logger.Error("faied to discover API groups", zap.Error(err)) + return false + } + return utils.SupportedMetricsAPIVersionAvailable(apiGroups) +} + func (gp *GenericPool) updateCPUUtilizationSvc() { - for { + var metricsApiAvailabe bool + checkDuration := 30 + + if !gp.checkMetricsApi() { + checkDuration = 180 + gp.logger.Error("Metrics API not available") + } + + serviceFunc := func() { podMetricsList, err := gp.metricsClient.MetricsV1beta1().PodMetricses(gp.namespace).List(context.TODO(), metav1.ListOptions{ LabelSelector: "managed=false", }) - if err != nil { gp.logger.Error("failed to fetch pod metrics list", zap.Error(err)) - } else { - gp.logger.Debug("pods found", zap.Any("length", len(podMetricsList.Items))) - for _, val := range podMetricsList.Items { - p, _ := resource.ParseQuantity("0m") - for _, container := range val.Containers { - p.Add(container.Usage["cpu"]) - } - if value, ok := gp.podFSVCMap.Load(val.ObjectMeta.Name); ok { - if valArray, ok1 := value.([]interface{}); ok1 { - function, address := valArray[0], valArray[1] - gp.fsCache.SetCPUUtilizaton(function.(string), address.(string), p) - gp.logger.Info(fmt.Sprintf("updated function %s, address %s, cpuUsage %+v", function.(string), address.(string), p)) - } + return + } + gp.logger.Debug("pods found", zap.Any("length", len(podMetricsList.Items))) + for _, val := range podMetricsList.Items { + p, _ := resource.ParseQuantity("0m") + for _, container := range val.Containers { + p.Add(container.Usage["cpu"]) + } + if value, ok := gp.podFSVCMap.Load(val.ObjectMeta.Name); ok { + if valArray, ok1 := value.([]interface{}); ok1 { + function, address := valArray[0], valArray[1] + gp.fsCache.SetCPUUtilizaton(function.(string), address.(string), p) + gp.logger.Info(fmt.Sprintf("updated function %s, address %s, cpuUsage %+v", function.(string), address.(string), p)) } } } + } - time.Sleep(30 * time.Second) + for { + if metricsApiAvailabe { + serviceFunc() + } else { + if gp.checkMetricsApi() { + metricsApiAvailabe = true + checkDuration = 30 + } + } + time.Sleep(time.Duration(checkDuration) * time.Second) } } diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index 913c7a02..9f98ffd8 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -211,7 +211,7 @@ func (gpm *GenericPoolManager) getPodInfo(obj apiv1.ObjectReference) (*apiv1.Pod } if err != nil || !exists { - gpm.logger.Debug("Falling back to getting pod info from k8s API -- this may cause performace issues for your function.") + gpm.logger.Debug("Falling back to getting pod info from k8s API -- this may cause performance issues for your function.") pod, err := gpm.kubernetesClient.CoreV1().Pods(obj.Namespace).Get(context.TODO(), obj.Name, metav1.GetOptions{}) return pod, err } diff --git a/pkg/executor/executortype/poolmgr/packagewatcher.go b/pkg/executor/executortype/poolmgr/packagewatcher.go index 917a716d..296304be 100644 --- a/pkg/executor/executortype/poolmgr/packagewatcher.go +++ b/pkg/executor/executortype/poolmgr/packagewatcher.go @@ -19,7 +19,6 @@ package poolmgr import ( "time" - "github.com/fission/fission/pkg/utils" "go.uber.org/zap" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/fields" @@ -28,6 +27,7 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" + "github.com/fission/fission/pkg/utils" ) // TODO : It may make sense to make each of add, update, delete funcs run as separate go routines. diff --git a/pkg/executor/fscache/metrics.go b/pkg/executor/fscache/metrics.go index 64f4049d..b3e58bc5 100644 --- a/pkg/executor/fscache/metrics.go +++ b/pkg/executor/fscache/metrics.go @@ -88,7 +88,7 @@ func (fsc *FunctionServiceCache) setFuncAlive(funcname, funcuid string, isAlive // ReapTime is the amount of time taken to reap a pod func (fsc *FunctionServiceCache) ReapTime(funcName, funcAddress string, time float64) { - funcReapTime.WithLabelValues(funcName, funcAddress).Observe(float64(time)) + funcReapTime.WithLabelValues(funcName, funcAddress).Observe(time) } // IdleTime is the amount of time it took Reaper to find out the pod was idle diff --git a/pkg/fission-cli/cmd/environment/create.go b/pkg/fission-cli/cmd/environment/create.go index 87536201..44f73bd7 100644 --- a/pkg/fission-cli/cmd/environment/create.go +++ b/pkg/fission-cli/cmd/environment/create.go @@ -129,7 +129,7 @@ func createEnvironmentFromCmd(input cli.Input) (*fv1.Environment, error) { poolsize := input.Int(flagkey.EnvPoolsize) if poolsize < 1 { - console.Warn("poolsize is not positive, if you are using pool manager please set postive value") + console.Warn("poolsize is not positive, if you are using pool manager please set positive value") } envBuilderImg := input.String(flagkey.EnvBuilderImage) diff --git a/pkg/fission-cli/cmd/environment/update.go b/pkg/fission-cli/cmd/environment/update.go index 55fc2e07..e113b28b 100644 --- a/pkg/fission-cli/cmd/environment/update.go +++ b/pkg/fission-cli/cmd/environment/update.go @@ -105,7 +105,7 @@ func updateExistingEnvironmentWithCmd(env *fv1.Environment, input cli.Input) (*f if input.IsSet(flagkey.EnvPoolsize) { env.Spec.Poolsize = input.Int(flagkey.EnvPoolsize) if env.Spec.Poolsize < 1 { - console.Warn("poolsize is not positive, if you are using pool manager please set postive value") + console.Warn("poolsize is not positive, if you are using pool manager please set positive value") } } diff --git a/pkg/fission-cli/cmd/package/package.go b/pkg/fission-cli/cmd/package/package.go index 3891fc64..48b5d754 100644 --- a/pkg/fission-cli/cmd/package/package.go +++ b/pkg/fission-cli/cmd/package/package.go @@ -28,6 +28,7 @@ import ( "github.com/hashicorp/go-multierror" "github.com/mholt/archiver" "github.com/pkg/errors" + uuid "github.com/satori/go.uuid" fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/controller/client" @@ -39,7 +40,6 @@ import ( flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" "github.com/fission/fission/pkg/fission-cli/util" "github.com/fission/fission/pkg/utils" - uuid "github.com/satori/go.uuid" ) // CreateArchive returns a fv1.Archive made from an archive . If specFile, then diff --git a/pkg/mqtrigger/messageQueue/kafka/kafka.go b/pkg/mqtrigger/messageQueue/kafka/kafka.go index 38332fe6..07653266 100644 --- a/pkg/mqtrigger/messageQueue/kafka/kafka.go +++ b/pkg/mqtrigger/messageQueue/kafka/kafka.go @@ -315,7 +315,7 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.Me zap.String("body", string(body))) if err != nil { - errorString := string("request body error: " + string(body)) + errorString := "request body error: " + string(body) errorHeaders := generateErrorHeaders(errorString) errorHandler(kafka.logger, trigger, producer, url, errors.Wrapf(err, errorString), errorHeaders) diff --git a/pkg/mqtrigger/scalermanager_test.go b/pkg/mqtrigger/scalermanager_test.go index ef06f597..9ec89b36 100644 --- a/pkg/mqtrigger/scalermanager_test.go +++ b/pkg/mqtrigger/scalermanager_test.go @@ -7,13 +7,14 @@ import ( "sort" "testing" - fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/stretchr/testify/assert" apiv1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/client-go/kubernetes" "k8s.io/client-go/kubernetes/fake" + + fv1 "github.com/fission/fission/pkg/apis/core/v1" ) func Test_toEnvVar(t *testing.T) { diff --git a/pkg/poolcache/poolcache.go b/pkg/poolcache/poolcache.go index 1382fdb9..9da8dafc 100644 --- a/pkg/poolcache/poolcache.go +++ b/pkg/poolcache/poolcache.go @@ -21,8 +21,9 @@ package poolcache import ( "fmt" - ferror "github.com/fission/fission/pkg/error" "k8s.io/apimachinery/pkg/api/resource" + + ferror "github.com/fission/fission/pkg/error" ) type requestType int diff --git a/pkg/router/functionHandler.go b/pkg/router/functionHandler.go index 10da3deb..8fb6bd59 100644 --- a/pkg/router/functionHandler.go +++ b/pkg/router/functionHandler.go @@ -534,7 +534,7 @@ func (roundTripper RetryingRoundTripper) addForwardedHostHeader(req *http.Reques // unTapservice marks the serviceURL in executor's cache as inactive, so that it can be reused func (fh functionHandler) unTapService(fn *fv1.Function, serviceUrl *url.URL) error { - fh.logger.Info("UnTapService Called") + fh.logger.Debug("UnTapService Called") ctx, cancel := context.WithTimeout(context.Background(), fh.unTapServiceTimeout) defer cancel() err := fh.executor.UnTapService(ctx, fn.ObjectMeta, fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType, serviceUrl) diff --git a/pkg/storagesvc/archivePruner.go b/pkg/storagesvc/archivePruner.go index a8586d32..87bcd505 100644 --- a/pkg/storagesvc/archivePruner.go +++ b/pkg/storagesvc/archivePruner.go @@ -76,7 +76,7 @@ func (pruner *ArchivePruner) insertArchive(archiveID string) { // and not the archives that are referenced by them, leaving the archives as orphans. // getOrphanArchives reaps the orphaned archives. func (pruner *ArchivePruner) getOrphanArchives() { - pruner.logger.Info("getting orphan archives") + pruner.logger.Debug("getting orphan archives") archivesRefByPkgs := make([]string, 0) var archiveID string diff --git a/pkg/storagesvc/client/storagesvc_test.go b/pkg/storagesvc/client/storagesvc_test.go index b81f37dd..c5dedbc9 100644 --- a/pkg/storagesvc/client/storagesvc_test.go +++ b/pkg/storagesvc/client/storagesvc_test.go @@ -28,12 +28,13 @@ import ( "time" "github.com/dchest/uniuri" - "github.com/fission/fission/pkg/storagesvc" "github.com/minio/minio-go" "github.com/ory/dockertest" dc "github.com/ory/dockertest/docker" "go.uber.org/zap" "go.uber.org/zap/zapcore" + + "github.com/fission/fission/pkg/storagesvc" ) const ( @@ -153,7 +154,7 @@ func TestS3StorageService(t *testing.T) { time.Sleep(10 * time.Second) - // Retrive file through minioClient + // Retrieve file through minioClient reader, err := minioClient.GetObject(bucketName, fileID, minio.GetObjectOptions{}) panicIf(err) defer reader.Close() diff --git a/pkg/storagesvc/stowClient.go b/pkg/storagesvc/stowClient.go index 782aa4b5..edd337e0 100644 --- a/pkg/storagesvc/stowClient.go +++ b/pkg/storagesvc/stowClient.go @@ -114,7 +114,7 @@ func (client *StowClient) putFile(file multipart.File, fileSize int64) (string, uploadName := client.config.storage.getUploadFileName() // save the file to the storage backend - item, err := client.container.Put(uploadName, file, int64(fileSize), nil) + item, err := client.container.Put(uploadName, file, fileSize, nil) if err != nil { client.logger.Error("error writing file on storage", zap.Error(err), diff --git a/pkg/utils/informer.go b/pkg/utils/informer.go index 094d4de0..8444fa47 100644 --- a/pkg/utils/informer.go +++ b/pkg/utils/informer.go @@ -1,12 +1,14 @@ package utils import ( - v1 "github.com/fission/fission/pkg/apis/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/selection" k8sInformers "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" + metricsapi "k8s.io/metrics/pkg/apis/metrics" + + v1 "github.com/fission/fission/pkg/apis/core/v1" ) func GetInformerFacoryByExecutor(client *kubernetes.Clientset, executorType v1.ExecutorType) (k8sInformers.SharedInformerFactory, error) { @@ -22,3 +24,22 @@ func GetInformerFacoryByExecutor(client *kubernetes.Clientset, executorType v1.E })) return informerFactory, nil } + +func SupportedMetricsAPIVersionAvailable(discoveredAPIGroups *metav1.APIGroupList) bool { + var supportedMetricsAPIVersions = []string{ + "v1beta1", + } + for _, discoveredAPIGroup := range discoveredAPIGroups.Groups { + if discoveredAPIGroup.Name != metricsapi.GroupName { + continue + } + for _, version := range discoveredAPIGroup.Versions { + for _, supportedVersion := range supportedMetricsAPIVersions { + if version.Version == supportedVersion { + return true + } + } + } + } + return false +}