diff --git a/cmd/fission-bundle/mqtrigger/mqtrigger.go b/cmd/fission-bundle/mqtrigger/mqtrigger.go index 7448eac0..6e4c5f77 100644 --- a/cmd/fission-bundle/mqtrigger/mqtrigger.go +++ b/cmd/fission-bundle/mqtrigger/mqtrigger.go @@ -35,10 +35,10 @@ import ( ) func Start(ctx context.Context, logger *zap.Logger, routerUrl string) error { - fissionClient, _, _, _, err := crd.MakeFissionClient() - + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - return errors.Wrap(err, "failed to get fission or kubernetes client") + return errors.Wrap(err, "failed to get fission client") } err = crd.WaitForCRDs(ctx, logger, fissionClient) diff --git a/cmd/preupgradechecks/checks.go b/cmd/preupgradechecks/checks.go index f01692fc..3304ff47 100644 --- a/cmd/preupgradechecks/checks.go +++ b/cmd/preupgradechecks/checks.go @@ -52,9 +52,18 @@ const ( ) func makePreUpgradeTaskClient(logger *zap.Logger) (*PreUpgradeTaskClient, error) { - fissionClient, k8sClient, apiExtClient, _, err := crd.MakeFissionClient() + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - return nil, errors.Wrap(err, "error making fission client") + return nil, errors.Wrap(err, "failed to get fission client") + } + k8sClient, err := clientGen.GetKubernetesClient() + if err != nil { + return nil, errors.Wrap(err, "failed to get kubernetes client") + } + apiExtClient, err := clientGen.GetApiExtensionsClient() + if err != nil { + return nil, errors.Wrap(err, "failed to get apiextensions client") } return &PreUpgradeTaskClient{ diff --git a/pkg/buildermgr/buildermgr.go b/pkg/buildermgr/buildermgr.go index 9fc2feea..8bcb8f29 100644 --- a/pkg/buildermgr/buildermgr.go +++ b/pkg/buildermgr/buildermgr.go @@ -34,9 +34,14 @@ import ( func Start(ctx context.Context, logger *zap.Logger, storageSvcUrl string) error { bmLogger := logger.Named("builder_manager") - fissionClient, kubernetesClient, _, _, err := crd.MakeFissionClient() + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - return errors.Wrap(err, "failed to get fission or kubernetes client") + return errors.Wrap(err, "failed to get fission client") + } + kubernetesClient, err := clientGen.GetKubernetesClient() + if err != nil { + return errors.Wrap(err, "failed to get kubernetes client") } err = crd.WaitForCRDs(ctx, logger, fissionClient) diff --git a/pkg/canaryconfigmgr/canaryConfigMgr.go b/pkg/canaryconfigmgr/canaryConfigMgr.go index a98adb9b..65d46d0a 100644 --- a/pkg/canaryconfigmgr/canaryConfigMgr.go +++ b/pkg/canaryconfigmgr/canaryConfigMgr.go @@ -555,12 +555,17 @@ func getEnvValue(envVar string) string { func StartCanaryServer(ctx context.Context, logger *zap.Logger, unitTestFlag bool) error { cLogger := logger.Named("CanaryServer") - fc, kc, _, _, err := crd.MakeFissionClient() + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - cLogger.Fatal("failed to connect to k8s API", zap.Error(err)) + return errors.Wrap(err, "failed to get fission client") + } + kubernetesClient, err := clientGen.GetKubernetesClient() + if err != nil { + return errors.Wrap(err, "failed to get kubernetes client") } - err = ConfigureFeatures(ctx, cLogger, unitTestFlag, fc, kc) + err = ConfigureFeatures(ctx, cLogger, unitTestFlag, fissionClient, kubernetesClient) if err != nil { cLogger.Error("error configuring features - proceeding without optional features", zap.Error(err)) } diff --git a/pkg/controller/api.go b/pkg/controller/api.go index 07f63384..54063cca 100644 --- a/pkg/controller/api.go +++ b/pkg/controller/api.go @@ -68,8 +68,8 @@ type ( } ) -func MakeAPI(logger *zap.Logger) (*API, error) { - api, err := makeCRDBackedAPI(logger) +func MakeAPI(logger *zap.Logger, fissionClient versioned.Interface, kubernetesClient kubernetes.Interface) (*API, error) { + api, err := makeCRDBackedAPI(logger, fissionClient, kubernetesClient) u := os.Getenv("STORAGE_SERVICE_URL") if len(u) > 0 { diff --git a/pkg/controller/api_test.go b/pkg/controller/api_test.go index e6d13b16..363b2670 100644 --- a/pkg/controller/api_test.go +++ b/pkg/controller/api_test.go @@ -503,7 +503,8 @@ func TestMain(m *testing.M) { return } - _, kubeClient, _, _, err := crd.GetKubernetesClient() + clientGen := crd.NewClientGenerator() + kubeClient, err := clientGen.GetKubernetesClient() panicIf(err) // testNS isolation for running multiple CI builds concurrently. diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 0f54cec9..9ec238bb 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -27,9 +27,18 @@ import ( func Start(ctx context.Context, logger *zap.Logger, port int, unitTestFlag bool) { cLogger := logger.Named("controller") - fc, _, apiExtClient, _, err := crd.MakeFissionClient() + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - cLogger.Fatal("failed to connect to k8s API", zap.Error(err)) + cLogger.Fatal("failed to get fission client", zap.Error(err)) + } + apiExtClient, err := clientGen.GetApiExtensionsClient() + if err != nil { + cLogger.Fatal("failed to get api extension client client", zap.Error(err)) + } + kubeClient, err := clientGen.GetKubernetesClient() + if err != nil { + cLogger.Fatal("failed to get kubernetes client", zap.Error(err)) } err = crd.EnsureFissionCRDs(ctx, cLogger, apiExtClient) @@ -37,12 +46,12 @@ func Start(ctx context.Context, logger *zap.Logger, port int, unitTestFlag bool) cLogger.Fatal("failed to find fission CRDs", zap.Error(err)) } - err = crd.WaitForCRDs(ctx, logger, fc) + err = crd.WaitForCRDs(ctx, logger, fissionClient) if err != nil { cLogger.Fatal("error waiting for CRDs", zap.Error(err)) } - api, err := MakeAPI(cLogger) + api, err := MakeAPI(cLogger, fissionClient, kubeClient) if err != nil { cLogger.Fatal("failed to start controller", zap.Error(err)) } diff --git a/pkg/controller/crd.go b/pkg/controller/crd.go index e807b0ff..ac7a48b2 100644 --- a/pkg/controller/crd.go +++ b/pkg/controller/crd.go @@ -18,15 +18,12 @@ package controller import ( "go.uber.org/zap" + "k8s.io/client-go/kubernetes" - "github.com/fission/fission/pkg/crd" + "github.com/fission/fission/pkg/generated/clientset/versioned" ) -func makeCRDBackedAPI(logger *zap.Logger) (*API, error) { - fissionClient, kubernetesClient, _, _, err := crd.MakeFissionClient() - if err != nil { - return nil, err - } +func makeCRDBackedAPI(logger *zap.Logger, fissionClient versioned.Interface, kubernetesClient kubernetes.Interface) (*API, error) { return &API{ logger: logger.Named("api"), fissionClient: fissionClient, diff --git a/pkg/crd/client.go b/pkg/crd/client.go index 080b9558..71e7ed56 100644 --- a/pkg/crd/client.go +++ b/pkg/crd/client.go @@ -19,7 +19,6 @@ package crd import ( "context" "errors" - "os" "time" "go.uber.org/zap" @@ -29,64 +28,76 @@ import ( "k8s.io/client-go/kubernetes" _ "k8s.io/client-go/plugin/pkg/client/auth" "k8s.io/client-go/rest" - "k8s.io/client-go/tools/clientcmd" metricsclient "k8s.io/metrics/pkg/client/clientset/versioned" + "sigs.k8s.io/controller-runtime/pkg/client/config" "github.com/fission/fission/pkg/generated/clientset/versioned" "github.com/fission/fission/pkg/utils" ) -// GetKubernetesClient gets a kubernetes client using the kubeconfig file at the -// environment var $KUBECONFIG, or an in-cluster config if that's -// undefined. -func GetKubernetesClient() (*rest.Config, kubernetes.Interface, apiextensionsclient.Interface, metricsclient.Interface, error) { - var config *rest.Config - var err error - - // get the config, either from kubeconfig or using our - // in-cluster service account - kubeConfig := os.Getenv("KUBECONFIG") - if len(kubeConfig) != 0 { - config, err = clientcmd.BuildConfigFromFlags("", kubeConfig) - if err != nil { - return nil, nil, nil, nil, err - } - } else { - config, err = rest.InClusterConfig() - if err != nil { - return nil, nil, nil, nil, err - } - } - - // creates the clientset - clientset, err := kubernetes.NewForConfig(config) - if err != nil { - return nil, nil, nil, nil, err - } - - apiExtClientset, err := apiextensionsclient.NewForConfig(config) - if err != nil { - return nil, nil, nil, nil, err - } - - metricsClient, _ := metricsclient.NewForConfig(config) - - return config, clientset, apiExtClientset, metricsClient, nil +type ClientGenerator struct { + restConfig *rest.Config } -func MakeFissionClient() (versioned.Interface, kubernetes.Interface, apiextensionsclient.Interface, metricsclient.Interface, error) { - config, kubeClient, apiExtClient, metricsClient, err := GetKubernetesClient() - if err != nil { - return nil, nil, nil, nil, err +func (cg *ClientGenerator) getRestConfig() (*rest.Config, error) { + if cg.restConfig != nil { + return cg.restConfig, nil } - // make a CRD REST client with the config - crdClient, err := versioned.NewForConfig(config) + var err error + cg.restConfig, err = config.GetConfig() if err != nil { - return nil, nil, nil, nil, err + return nil, err } + return cg.restConfig, nil +} - return crdClient, kubeClient, apiExtClient, metricsClient, nil +func (cg *ClientGenerator) GetFissionClient() (versioned.Interface, error) { + config, err := cg.getRestConfig() + if err != nil { + return nil, err + } + return versioned.NewForConfig(config) +} + +func (cg *ClientGenerator) GetKubernetesClient() (kubernetes.Interface, error) { + config, err := cg.getRestConfig() + if err != nil { + return nil, err + } + return kubernetes.NewForConfig(config) +} + +func (cg *ClientGenerator) GetApiExtensionsClient() (apiextensionsclient.Interface, error) { + config, err := cg.getRestConfig() + if err != nil { + return nil, err + } + return apiextensionsclient.NewForConfig(config) +} + +func (cg *ClientGenerator) GetMetricsClient() (metricsclient.Interface, error) { + config, err := cg.getRestConfig() + if err != nil { + return nil, err + } + return metricsclient.NewForConfig(config) +} + +func (cg *ClientGenerator) GetDynamicClient() (dynamic.Interface, error) { + config, err := cg.getRestConfig() + if err != nil { + return nil, err + } + return dynamic.NewForConfig(config) +} + +func NewClientGenerator() *ClientGenerator { + return &ClientGenerator{} +} + +func NewClientGeneratorWithRestConfig(restConfig *rest.Config) *ClientGenerator { + return &ClientGenerator{restConfig: restConfig} } // WaitForCRDs does a timeout to check if CRDs have been installed @@ -111,29 +122,3 @@ func WaitForCRDs(ctx context.Context, logger *zap.Logger, fissionClient versione } } } - -// GetDynamicClient creates and returns new dynamic client or returns an error -func GetDynamicClient() (dynamic.Interface, error) { - var config *rest.Config - var err error - - // get the config, either from kubeconfig or using our - // in-cluster service account - kubeConfig := os.Getenv("KUBECONFIG") - if len(kubeConfig) != 0 { - config, err = clientcmd.BuildConfigFromFlags("", kubeConfig) - if err != nil { - return nil, err - } - } else { - config, err = rest.InClusterConfig() - if err != nil { - return nil, err - } - } - dynamicClient, err := dynamic.NewForConfig(config) - if err != nil { - return nil, err - } - return dynamicClient, nil -} diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index 733c1447..2e913de7 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -251,9 +251,18 @@ func (executor *Executor) getFunctionServiceFromCache(ctx context.Context, fn *f // StartExecutor Starts executor and the executor components such as Poolmgr, // deploymgr and potential future executor types func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error { - fissionClient, kubernetesClient, _, metricsClient, err := crd.MakeFissionClient() + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - return errors.Wrap(err, "failed to get kubernetes client") + return errors.Wrap(err, "error making the fission client") + } + kubernetesClient, err := clientGen.GetKubernetesClient() + if err != nil { + return errors.Wrap(err, "error making the kube client") + } + metricsClient, err := clientGen.GetMetricsClient() + if err != nil { + logger.Error("error making the metrics client", zap.Error(err)) } err = crd.WaitForCRDs(ctx, logger, fissionClient) diff --git a/pkg/executor/executor_test.go b/pkg/executor/executor_test.go index 8af11808..ef92491e 100644 --- a/pkg/executor/executor_test.go +++ b/pkg/executor/executor_test.go @@ -115,9 +115,18 @@ func TestExecutor(t *testing.T) { // connect to k8s // and get CRD client - fissionClient, kubeClient, apiExtClient, _, err := crd.MakeFissionClient() + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - log.Panicf("failed to connect: %v", err) + log.Panicf("failed to connect: %s", err) + } + kubeClient, err := clientGen.GetKubernetesClient() + if err != nil { + log.Panicf("failed to connect: %s", err) + } + apiExtClient, err := clientGen.GetApiExtensionsClient() + if err != nil { + log.Panicf("failed to connect: %s", err) } ctx := context.Background() diff --git a/pkg/executor/metrics/metrics.go b/pkg/executor/metrics/metrics.go index 701cf38d..6a149906 100644 --- a/pkg/executor/metrics/metrics.go +++ b/pkg/executor/metrics/metrics.go @@ -17,8 +17,8 @@ limitations under the License. package metrics import ( + "github.com/fission/fission/pkg/utils/metrics" "github.com/prometheus/client_golang/prometheus" - "github.com/prometheus/client_golang/prometheus/promauto" ) var ( @@ -26,14 +26,14 @@ var ( // function_uid: the function's version id // function_address: the address of the pod from which the function was called functionLabels = []string{"function_name", "function_namespace"} - ColdStarts = promauto.NewCounterVec( + ColdStarts = prometheus.NewCounterVec( prometheus.CounterOpts{ Name: "fission_function_cold_starts_total", Help: "How many cold starts are made by function_name, function_namespace.", }, functionLabels, ) - FuncRunningSummary = promauto.NewSummaryVec( + FuncRunningSummary = prometheus.NewSummaryVec( prometheus.SummaryOpts{ Name: "fission_function_running_seconds", Help: "The running time (last access - create) in seconds of the function.", @@ -41,7 +41,7 @@ var ( }, functionLabels, ) - FuncError = promauto.NewCounterVec( + FuncError = prometheus.NewCounterVec( prometheus.CounterOpts{ Name: "fission_function_cold_start_errors_total", Help: "Count of fission cold start errors", @@ -49,3 +49,10 @@ var ( functionLabels, ) ) + +func init() { + registry := metrics.Registry + registry.MustRegister(ColdStarts) + registry.MustRegister(FuncRunningSummary) + registry.MustRegister(FuncError) +} diff --git a/pkg/fetcher/fetcher.go b/pkg/fetcher/fetcher.go index 8d605797..8d51c436 100644 --- a/pkg/fetcher/fetcher.go +++ b/pkg/fetcher/fetcher.go @@ -90,9 +90,14 @@ func MakeFetcher(logger *zap.Logger, sharedVolumePath string, sharedSecretPath s fLogger.Fatal("error creating shared config directory", zap.Error(err), zap.String("directory", sharedConfigPath)) } - fissionClient, kubeClient, _, _, err := crd.MakeFissionClient() + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - return nil, errors.Wrap(err, "error making the fission / kube client") + return nil, errors.Wrap(err, "error making the fission client") + } + kubeClient, err := clientGen.GetKubernetesClient() + if err != nil { + return nil, errors.Wrap(err, "error making the kube client") } name, err := os.ReadFile(fv1.PodInfoMount + "/name") diff --git a/pkg/fission-cli/cmd/client.go b/pkg/fission-cli/cmd/client.go index b00dd662..d68c7c52 100644 --- a/pkg/fission-cli/cmd/client.go +++ b/pkg/fission-cli/cmd/client.go @@ -27,6 +27,7 @@ import ( "k8s.io/client-go/rest" "k8s.io/client-go/tools/clientcmd" + "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/fission-cli/console" "github.com/fission/fission/pkg/generated/clientset/versioned" ) @@ -118,13 +119,15 @@ func NewClient(opts ClientOptions) (*Client, error) { return nil, err } client.RestConfig = restConfig - clientset, err := kubernetes.NewForConfig(restConfig) + + clientGen := crd.NewClientGeneratorWithRestConfig(restConfig) + clientset, err := clientGen.GetKubernetesClient() if err != nil { return nil, err } client.KubernetesClient = clientset - fissionClientset, err := versioned.NewForConfig(restConfig) + fissionClientset, err := clientGen.GetFissionClient() if err != nil { return nil, err } diff --git a/pkg/kubewatcher/main.go b/pkg/kubewatcher/main.go index 4d876dc6..5860d64d 100644 --- a/pkg/kubewatcher/main.go +++ b/pkg/kubewatcher/main.go @@ -27,9 +27,14 @@ import ( ) func Start(ctx context.Context, logger *zap.Logger, routerUrl string) error { - fissionClient, kubeClient, _, _, err := crd.MakeFissionClient() + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - return errors.Wrap(err, "failed to get fission or kubernetes client") + return errors.Wrap(err, "failed to get fission client") + } + kubeClient, err := clientGen.GetKubernetesClient() + if err != nil { + return errors.Wrap(err, "failed to get kubernetes client") } err = crd.WaitForCRDs(ctx, logger, fissionClient) diff --git a/pkg/logger/logger.go b/pkg/logger/logger.go index 5cd22fea..a25f1e7b 100644 --- a/pkg/logger/logger.go +++ b/pkg/logger/logger.go @@ -167,7 +167,9 @@ func Start(ctx context.Context, logger *zap.Logger) { } } go symlinkReaper(logger) - _, kubernetesClient, _, _, err := crd.MakeFissionClient() + + clientGen := crd.NewClientGenerator() + kubernetesClient, err := clientGen.GetKubernetesClient() if err != nil { log.Fatalf("Error starting pod watcher: %v", err) } diff --git a/pkg/mqtrigger/metrics.go b/pkg/mqtrigger/metrics.go index 8ded7a3a..13d14fb0 100644 --- a/pkg/mqtrigger/metrics.go +++ b/pkg/mqtrigger/metrics.go @@ -18,26 +18,27 @@ package mqtrigger import ( "github.com/prometheus/client_golang/prometheus" - "github.com/prometheus/client_golang/prometheus/promauto" + + "github.com/fission/fission/pkg/utils/metrics" ) var ( labels = []string{"trigger_name", "trigger_namespace"} - subscriptionCount = promauto.NewGaugeVec( + subscriptionCount = prometheus.NewGaugeVec( prometheus.GaugeOpts{ Name: "fission_mqt_subscriptions", Help: "Total number of subscriptions to mq currently", }, []string{}, ) - messageCount = promauto.NewCounterVec( + messageCount = prometheus.NewCounterVec( prometheus.CounterOpts{ Name: "fission_mqt_messages_processed_total", Help: "Total number of messages processed", }, labels, ) - messageLagCount = promauto.NewGaugeVec( + messageLagCount = prometheus.NewGaugeVec( prometheus.GaugeOpts{ Name: "fission_mqt_message_lag", Help: "Total number of messages lag per topic and partition", @@ -61,3 +62,10 @@ func IncreaseMessageCount(trigname, trignamespace string) { func SetMessageLagCount(trigname, trignamespace, topic, partition string, lag int64) { messageLagCount.WithLabelValues(trigname, trignamespace, topic, partition).Set(float64(lag)) } + +func init() { + registry := metrics.Registry + registry.MustRegister(subscriptionCount) + registry.MustRegister(messageCount) + registry.MustRegister(messageLagCount) +} diff --git a/pkg/mqtrigger/scalermanager.go b/pkg/mqtrigger/scalermanager.go index c83062e8..a9ffc967 100644 --- a/pkg/mqtrigger/scalermanager.go +++ b/pkg/mqtrigger/scalermanager.go @@ -50,7 +50,8 @@ var ( ) func getScaledObjectClient(namespace string) (dynamic.ResourceInterface, error) { - dynamicClient, err := crd.GetDynamicClient() + clientGen := crd.NewClientGenerator() + dynamicClient, err := clientGen.GetDynamicClient() if err != nil { return nil, err } @@ -58,7 +59,8 @@ func getScaledObjectClient(namespace string) (dynamic.ResourceInterface, error) } func getAuthTriggerClient(namespace string) (dynamic.ResourceInterface, error) { - dynamicClient, err := crd.GetDynamicClient() + clientGen := crd.NewClientGenerator() + dynamicClient, err := clientGen.GetDynamicClient() if err != nil { return nil, err } @@ -149,10 +151,16 @@ func mqTriggerEventHandlers(ctx context.Context, logger *zap.Logger, kubeClient // StartScalerManager watches for changes in MessageQueueTrigger and, // Based on changes, it Creates, Updates and Deletes Objects of Kind ScaledObjects, AuthenticationTriggers and Deployments func StartScalerManager(ctx context.Context, logger *zap.Logger, routerURL string) error { - fissionClient, kubeClient, _, _, err := crd.MakeFissionClient() + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - return err + return errors.Wrap(err, "failed to get fission client") } + kubeClient, err := clientGen.GetKubernetesClient() + if err != nil { + return errors.Wrap(err, "failed to get kubernetes client") + } + err = crd.WaitForCRDs(ctx, logger, fissionClient) if err != nil { return errors.Wrap(err, "error waiting for CRDs") diff --git a/pkg/router/metrics.go b/pkg/router/metrics.go index fb08e424..caf34f4a 100644 --- a/pkg/router/metrics.go +++ b/pkg/router/metrics.go @@ -2,7 +2,8 @@ package router import ( "github.com/prometheus/client_golang/prometheus" - "github.com/prometheus/client_golang/prometheus/promauto" + + "github.com/fission/fission/pkg/utils/metrics" ) var ( @@ -15,21 +16,21 @@ var ( // code: http status code // path: the client call the function on which http path // method: the function's http method - functionCalls = promauto.NewCounterVec( + functionCalls = prometheus.NewCounterVec( prometheus.CounterOpts{ Name: "fission_function_calls_total", Help: "Count of Fission function calls", }, labelsStrings, ) - functionCallErrors = promauto.NewCounterVec( + functionCallErrors = prometheus.NewCounterVec( prometheus.CounterOpts{ Name: "fission_function_errors_total", Help: "Count of Fission function errors", }, labelsStrings, ) - functionCallOverhead = promauto.NewSummaryVec( + functionCallOverhead = prometheus.NewSummaryVec( prometheus.SummaryOpts{ Name: "fission_function_overhead_seconds", Help: "The function call delay caused by fission.", @@ -38,3 +39,10 @@ var ( labelsStrings, ) ) + +func init() { + registry := metrics.Registry + registry.MustRegister(functionCalls) + registry.MustRegister(functionCallErrors) + registry.MustRegister(functionCallOverhead) +} diff --git a/pkg/router/router.go b/pkg/router/router.go index 55608977..199f8767 100644 --- a/pkg/router/router.go +++ b/pkg/router/router.go @@ -90,9 +90,14 @@ func serve(ctx context.Context, logger *zap.Logger, port int, func Start(ctx context.Context, logger *zap.Logger, port int, executorURL string) { fmap := makeFunctionServiceMap(logger, time.Minute) - fissionClient, kubeClient, _, _, err := crd.MakeFissionClient() + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - logger.Fatal("error connecting to kubernetes API", zap.Error(err)) + logger.Fatal("error making the fission client", zap.Error(err)) + } + kubeClient, err := clientGen.GetKubernetesClient() + if err != nil { + logger.Fatal("error making the kube client", zap.Error(err)) } err = crd.WaitForCRDs(ctx, logger, fissionClient) diff --git a/pkg/storagesvc/archivePruner.go b/pkg/storagesvc/archivePruner.go index 1c450c99..9f73324a 100644 --- a/pkg/storagesvc/archivePruner.go +++ b/pkg/storagesvc/archivePruner.go @@ -26,6 +26,7 @@ import ( "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/generated/clientset/versioned" "github.com/fission/fission/pkg/utils" + "github.com/pkg/errors" ) type ArchivePruner struct { @@ -39,14 +40,15 @@ type ArchivePruner struct { const defaultPruneInterval int = 60 // in minutes func MakeArchivePruner(logger *zap.Logger, stowClient *StowClient, pruneInterval time.Duration) (*ArchivePruner, error) { - crdClient, _, _, _, err := crd.MakeFissionClient() + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - return nil, err + return nil, errors.Wrap(err, "failed to get fission client") } return &ArchivePruner{ logger: logger.Named("archive_pruner"), - crdClient: crdClient, + crdClient: fissionClient, archiveChan: make(chan string), stowClient: stowClient, pruneInterval: pruneInterval, diff --git a/pkg/storagesvc/metrics.go b/pkg/storagesvc/metrics.go index ed60e71d..0af8d1d4 100644 --- a/pkg/storagesvc/metrics.go +++ b/pkg/storagesvc/metrics.go @@ -2,19 +2,20 @@ package storagesvc import ( "github.com/prometheus/client_golang/prometheus" - "github.com/prometheus/client_golang/prometheus/promauto" + + "github.com/fission/fission/pkg/utils/metrics" ) var ( functionLabels = []string{} - totalArchives = promauto.NewGaugeVec( + totalArchives = prometheus.NewGaugeVec( prometheus.GaugeOpts{ Name: "fission_archives", Help: "Number of archives stored", }, functionLabels, ) - totalMemoryUsage = promauto.NewGaugeVec( + totalMemoryUsage = prometheus.NewGaugeVec( prometheus.GaugeOpts{ Name: "fission_archive_memory_bytes", Help: "Amount of memory consumed by archives", @@ -22,3 +23,9 @@ var ( functionLabels, ) ) + +func init() { + registry := metrics.Registry + registry.MustRegister(totalArchives) + registry.MustRegister(totalMemoryUsage) +} diff --git a/pkg/timer/main.go b/pkg/timer/main.go index 4383dfbb..81a327d6 100644 --- a/pkg/timer/main.go +++ b/pkg/timer/main.go @@ -27,9 +27,10 @@ import ( ) func Start(ctx context.Context, logger *zap.Logger, routerUrl string) error { - fissionClient, _, _, _, err := crd.MakeFissionClient() + clientGen := crd.NewClientGenerator() + fissionClient, err := clientGen.GetFissionClient() if err != nil { - return errors.Wrap(err, "failed to get fission or kubernetes client") + return errors.Wrap(err, "failed to get fission client") } err = crd.WaitForCRDs(ctx, logger, fissionClient) diff --git a/pkg/utils/metrics/http_metrics.go b/pkg/utils/metrics/http_metrics.go index 1a85a3d3..02549750 100644 --- a/pkg/utils/metrics/http_metrics.go +++ b/pkg/utils/metrics/http_metrics.go @@ -22,7 +22,6 @@ import ( "github.com/gorilla/mux" "github.com/prometheus/client_golang/prometheus" - "github.com/prometheus/client_golang/prometheus/promauto" "github.com/fission/fission/pkg/router/util" ) @@ -33,14 +32,14 @@ type ResponseWriterWrapper struct { } var ( - httpRequestsTotal = promauto.NewCounterVec( + httpRequestsTotal = prometheus.NewCounterVec( prometheus.CounterOpts{ Name: "http_requests_total", Help: "Number of requests by path, method and status code.", }, []string{"path", "method", "code"}, ) - httpRequestDuration = promauto.NewSummaryVec( + httpRequestDuration = prometheus.NewSummaryVec( prometheus.SummaryOpts{ Name: "http_requests_duration_seconds", Help: "Time taken to serve the request by path and method.", @@ -48,7 +47,7 @@ var ( }, []string{"path", "method"}, ) - httpRequestInFlight = promauto.NewGaugeVec( + httpRequestInFlight = prometheus.NewGaugeVec( prometheus.GaugeOpts{ Name: "http_requests_in_flight", Help: "Number of requests currently being served by path and method.", @@ -57,6 +56,12 @@ var ( ) ) +func init() { + Registry.MustRegister(httpRequestsTotal) + Registry.MustRegister(httpRequestDuration) + Registry.MustRegister(httpRequestInFlight) +} + func (rw *ResponseWriterWrapper) WriteHeader(statuscode int) { rw.statusCode = statuscode rw.ResponseWriter.WriteHeader(statuscode) diff --git a/pkg/utils/metrics/registry.go b/pkg/utils/metrics/registry.go new file mode 100644 index 00000000..9793564c --- /dev/null +++ b/pkg/utils/metrics/registry.go @@ -0,0 +1,9 @@ +package metrics + +import ( + "github.com/prometheus/client_golang/prometheus" +) + +var ( + Registry = prometheus.NewRegistry() +) diff --git a/pkg/utils/metrics/server.go b/pkg/utils/metrics/server.go index 642b7390..ad2d9a5a 100644 --- a/pkg/utils/metrics/server.go +++ b/pkg/utils/metrics/server.go @@ -23,6 +23,7 @@ import ( "github.com/prometheus/client_golang/prometheus/promhttp" "go.uber.org/zap" + "sigs.k8s.io/controller-runtime/pkg/metrics" "github.com/fission/fission/pkg/utils/httpserver" ) @@ -32,7 +33,17 @@ func ServeMetrics(ctx context.Context, logger *zap.Logger) { if metricsAddr == "" { metricsAddr = "8080" } + err := metrics.Registry.Register(Registry) + if err != nil { + logger.Error("failed to register metrics", zap.Error(err)) + } mux := http.NewServeMux() - mux.Handle("/metrics", promhttp.Handler()) + mux.Handle("/metrics", promhttp.HandlerFor( + metrics.Registry, + promhttp.HandlerOpts{ + // Opt into OpenMetrics to support exemplars. + EnableOpenMetrics: true, + }, + )) httpserver.StartServer(ctx, logger, "metrics", metricsAddr, mux) }