Use client generator to generate all k8s clients and add respective client-go metrics (#2668)
* Define client generator to generate all k8s clients * Increase QPS and burst values * Capture client-go metrics * Support for controller runtime metrics Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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{
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
+57
-72
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
|
||||
+12
-4
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
+3
-2
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -0,0 +1,9 @@
|
||||
package metrics
|
||||
|
||||
import (
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
)
|
||||
|
||||
var (
|
||||
Registry = prometheus.NewRegistry()
|
||||
)
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user