From 154fe0d447b3d14c56e27743e4974d1a53c80d6a Mon Sep 17 00:00:00 2001 From: Sanket Sudake Date: Mon, 28 Jun 2021 10:06:32 +0530 Subject: [PATCH] Retrieve pod metrics only if metrics server is running and Go lint fixes (#2094) * Retrieve pod metrics only if metrics server is running Currently we query pod metrics every 30 sec which floods executor logs, added check which confirms if metrics server is running then only we start querying pod metrics for identifying CPU utilization. Signed-off-by: Sanket Sudake * Fixed couple of typos and misspells with Go CI Signed-off-by: Sanket Sudake * Remove unnecessary conversions with Go CI Signed-off-by: Sanket Sudake --- .golangci.yaml | 17 +++++- cmd/reporter/app/cmd_event.go | 3 +- pkg/builder/client/client.go | 2 +- pkg/controller/packageApi.go | 2 +- .../executortype/newdeploy/newdeploy.go | 2 +- .../executortype/newdeploy/newdeploymgr.go | 4 +- pkg/executor/executortype/poolmgr/gp.go | 60 +++++++++++++------ pkg/executor/executortype/poolmgr/gpm.go | 2 +- .../executortype/poolmgr/packagewatcher.go | 2 +- pkg/executor/fscache/metrics.go | 2 +- pkg/fission-cli/cmd/environment/create.go | 2 +- pkg/fission-cli/cmd/environment/update.go | 2 +- pkg/fission-cli/cmd/package/package.go | 2 +- pkg/mqtrigger/messageQueue/kafka/kafka.go | 2 +- pkg/mqtrigger/scalermanager_test.go | 3 +- pkg/poolcache/poolcache.go | 3 +- pkg/router/functionHandler.go | 2 +- pkg/storagesvc/archivePruner.go | 2 +- pkg/storagesvc/client/storagesvc_test.go | 5 +- pkg/storagesvc/stowClient.go | 2 +- pkg/utils/informer.go | 23 ++++++- 21 files changed, 105 insertions(+), 39 deletions(-) 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 +}