Add typed informers instead of generic SharedIndexInformers (#2174)
Generally using typed informers is more standard practise than using SharedIndexInformer(SII) implicity. SII also lack listers provided by informer factory and few other high level abstractions. Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
@@ -21,8 +21,8 @@ import (
|
|||||||
|
|
||||||
"github.com/pkg/errors"
|
"github.com/pkg/errors"
|
||||||
"go.uber.org/zap"
|
"go.uber.org/zap"
|
||||||
|
informerv1 "k8s.io/client-go/informers/core/v1"
|
||||||
"k8s.io/client-go/kubernetes"
|
"k8s.io/client-go/kubernetes"
|
||||||
k8sCache "k8s.io/client-go/tools/cache"
|
|
||||||
|
|
||||||
fv1 "github.com/fission/fission/pkg/apis/core/v1"
|
fv1 "github.com/fission/fission/pkg/apis/core/v1"
|
||||||
"github.com/fission/fission/pkg/crd"
|
"github.com/fission/fission/pkg/crd"
|
||||||
@@ -34,9 +34,6 @@ type (
|
|||||||
ConfigSecretController struct {
|
ConfigSecretController struct {
|
||||||
logger *zap.Logger
|
logger *zap.Logger
|
||||||
|
|
||||||
configmapInformer *k8sCache.SharedIndexInformer
|
|
||||||
secretInformer *k8sCache.SharedIndexInformer
|
|
||||||
|
|
||||||
fissionClient *crd.FissionClient
|
fissionClient *crd.FissionClient
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
@@ -44,17 +41,15 @@ type (
|
|||||||
// MakeConfigSecretController makes a controller for configmaps and secrets which changes related functions
|
// MakeConfigSecretController makes a controller for configmaps and secrets which changes related functions
|
||||||
func MakeConfigSecretController(ctx context.Context, logger *zap.Logger, fissionClient *crd.FissionClient,
|
func MakeConfigSecretController(ctx context.Context, logger *zap.Logger, fissionClient *crd.FissionClient,
|
||||||
kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType,
|
kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType,
|
||||||
configmapInformer *k8sCache.SharedIndexInformer,
|
configmapInformer informerv1.ConfigMapInformer,
|
||||||
secretInformer *k8sCache.SharedIndexInformer) *ConfigSecretController {
|
secretInformer informerv1.SecretInformer) *ConfigSecretController {
|
||||||
logger.Debug("Creating ConfigMap & Secret Controller")
|
logger.Debug("Creating ConfigMap & Secret Controller")
|
||||||
cmsController := &ConfigSecretController{
|
cmsController := &ConfigSecretController{
|
||||||
logger: logger,
|
logger: logger,
|
||||||
configmapInformer: configmapInformer,
|
fissionClient: fissionClient,
|
||||||
secretInformer: secretInformer,
|
|
||||||
fissionClient: fissionClient,
|
|
||||||
}
|
}
|
||||||
(*configmapInformer).AddEventHandler(ConfigMapEventHandlers(ctx, logger, fissionClient, kubernetesClient, types))
|
configmapInformer.Informer().AddEventHandler(ConfigMapEventHandlers(ctx, logger, fissionClient, kubernetesClient, types))
|
||||||
(*secretInformer).AddEventHandler(SecretEventHandlers(ctx, logger, fissionClient, kubernetesClient, types))
|
secretInformer.Informer().AddEventHandler(SecretEventHandlers(ctx, logger, fissionClient, kubernetesClient, types))
|
||||||
|
|
||||||
return cmsController
|
return cmsController
|
||||||
}
|
}
|
||||||
|
|||||||
+10
-10
@@ -277,15 +277,15 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames
|
|||||||
logger.Info("Starting executor", zap.String("instanceID", executorInstanceID))
|
logger.Info("Starting executor", zap.String("instanceID", executorInstanceID))
|
||||||
|
|
||||||
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30)
|
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30)
|
||||||
funcInformer := informerFactory.Core().V1().Functions().Informer()
|
funcInformer := informerFactory.Core().V1().Functions()
|
||||||
pkgInformer := informerFactory.Core().V1().Packages().Informer()
|
pkgInformer := informerFactory.Core().V1().Packages()
|
||||||
envInformer := informerFactory.Core().V1().Environments().Informer()
|
envInformer := informerFactory.Core().V1().Environments()
|
||||||
|
|
||||||
gpm, err := poolmgr.MakeGenericPoolManager(
|
gpm, err := poolmgr.MakeGenericPoolManager(
|
||||||
logger,
|
logger,
|
||||||
fissionClient, kubernetesClient, metricsClient,
|
fissionClient, kubernetesClient, metricsClient,
|
||||||
functionNamespace, fetcherConfig, executorInstanceID,
|
functionNamespace, fetcherConfig, executorInstanceID,
|
||||||
&funcInformer, &pkgInformer,
|
funcInformer, pkgInformer,
|
||||||
)
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return errors.Wrap(err, "pool manager creation faied")
|
return errors.Wrap(err, "pool manager creation faied")
|
||||||
@@ -295,7 +295,7 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames
|
|||||||
logger,
|
logger,
|
||||||
fissionClient, kubernetesClient,
|
fissionClient, kubernetesClient,
|
||||||
functionNamespace, fetcherConfig, executorInstanceID,
|
functionNamespace, fetcherConfig, executorInstanceID,
|
||||||
&funcInformer, &envInformer,
|
funcInformer, envInformer,
|
||||||
)
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return errors.Wrap(err, "new deploy manager creation faied")
|
return errors.Wrap(err, "new deploy manager creation faied")
|
||||||
@@ -306,7 +306,7 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames
|
|||||||
ctx,
|
ctx,
|
||||||
logger,
|
logger,
|
||||||
fissionClient, kubernetesClient,
|
fissionClient, kubernetesClient,
|
||||||
functionNamespace, executorInstanceID, &funcInformer)
|
functionNamespace, executorInstanceID, funcInformer)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return errors.Wrap(err, "container manager creation faied")
|
return errors.Wrap(err, "container manager creation faied")
|
||||||
}
|
}
|
||||||
@@ -334,13 +334,13 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames
|
|||||||
util.WaitTimeout(wg, 30*time.Second)
|
util.WaitTimeout(wg, 30*time.Second)
|
||||||
|
|
||||||
k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Minute*30)
|
k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Minute*30)
|
||||||
configmapInformer := k8sInformerFactory.Core().V1().ConfigMaps().Informer()
|
configmapInformer := k8sInformerFactory.Core().V1().ConfigMaps()
|
||||||
secretInformer := k8sInformerFactory.Core().V1().Secrets().Informer()
|
secretInformer := k8sInformerFactory.Core().V1().Secrets()
|
||||||
|
|
||||||
cms := cms.MakeConfigSecretController(ctx, logger, fissionClient, kubernetesClient, executorTypes, &configmapInformer, &secretInformer)
|
cms := cms.MakeConfigSecretController(ctx, logger, fissionClient, kubernetesClient, executorTypes, configmapInformer, secretInformer)
|
||||||
|
|
||||||
api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes, []k8sCache.SharedIndexInformer{
|
api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes, []k8sCache.SharedIndexInformer{
|
||||||
funcInformer, pkgInformer, envInformer, configmapInformer, secretInformer,
|
funcInformer.Informer(), pkgInformer.Informer(), envInformer.Informer(), configmapInformer.Informer(), secretInformer.Informer(),
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
|
|||||||
@@ -43,6 +43,7 @@ import (
|
|||||||
"github.com/fission/fission/pkg/executor/executortype"
|
"github.com/fission/fission/pkg/executor/executortype"
|
||||||
"github.com/fission/fission/pkg/executor/fscache"
|
"github.com/fission/fission/pkg/executor/fscache"
|
||||||
"github.com/fission/fission/pkg/executor/reaper"
|
"github.com/fission/fission/pkg/executor/reaper"
|
||||||
|
finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1"
|
||||||
"github.com/fission/fission/pkg/throttler"
|
"github.com/fission/fission/pkg/throttler"
|
||||||
"github.com/fission/fission/pkg/utils"
|
"github.com/fission/fission/pkg/utils"
|
||||||
"github.com/fission/fission/pkg/utils/maps"
|
"github.com/fission/fission/pkg/utils/maps"
|
||||||
@@ -68,7 +69,6 @@ type (
|
|||||||
|
|
||||||
throttler *throttler.Throttler
|
throttler *throttler.Throttler
|
||||||
|
|
||||||
funcInformer *k8sCache.SharedIndexInformer
|
|
||||||
serviceInformer k8sCache.SharedIndexInformer
|
serviceInformer k8sCache.SharedIndexInformer
|
||||||
deploymentInformer k8sCache.SharedIndexInformer
|
deploymentInformer k8sCache.SharedIndexInformer
|
||||||
|
|
||||||
@@ -84,7 +84,7 @@ func MakeContainer(
|
|||||||
kubernetesClient *kubernetes.Clientset,
|
kubernetesClient *kubernetes.Clientset,
|
||||||
namespace string,
|
namespace string,
|
||||||
instanceID string,
|
instanceID string,
|
||||||
funcInformer *k8sCache.SharedIndexInformer) (executortype.ExecutorType, error) {
|
funcInformer finformerv1.FunctionInformer) (executortype.ExecutorType, error) {
|
||||||
enableIstio := false
|
enableIstio := false
|
||||||
if len(os.Getenv("ENABLE_ISTIO")) > 0 {
|
if len(os.Getenv("ENABLE_ISTIO")) > 0 {
|
||||||
istio, err := strconv.ParseBool(os.Getenv("ENABLE_ISTIO"))
|
istio, err := strconv.ParseBool(os.Getenv("ENABLE_ISTIO"))
|
||||||
@@ -105,15 +105,13 @@ func MakeContainer(
|
|||||||
fsCache: fscache.MakeFunctionServiceCache(logger),
|
fsCache: fscache.MakeFunctionServiceCache(logger),
|
||||||
throttler: throttler.MakeThrottler(1 * time.Minute),
|
throttler: throttler.MakeThrottler(1 * time.Minute),
|
||||||
|
|
||||||
funcInformer: funcInformer,
|
|
||||||
|
|
||||||
runtimeImagePullPolicy: utils.GetImagePullPolicy(os.Getenv("RUNTIME_IMAGE_PULL_POLICY")),
|
runtimeImagePullPolicy: utils.GetImagePullPolicy(os.Getenv("RUNTIME_IMAGE_PULL_POLICY")),
|
||||||
useIstio: enableIstio,
|
useIstio: enableIstio,
|
||||||
// Time is set slightly higher than NewDeploy as cold starts are longer for CaaF
|
// Time is set slightly higher than NewDeploy as cold starts are longer for CaaF
|
||||||
defaultIdlePodReapTime: 1 * time.Minute,
|
defaultIdlePodReapTime: 1 * time.Minute,
|
||||||
}
|
}
|
||||||
|
|
||||||
(*caaf.funcInformer).AddEventHandler(caaf.FuncInformerHandler(ctx))
|
funcInformer.Informer().AddEventHandler(caaf.FuncInformerHandler(ctx))
|
||||||
|
|
||||||
informerFactory, err := utils.GetInformerFactoryByExecutor(caaf.kubernetesClient, fv1.ExecutorTypeContainer)
|
informerFactory, err := utils.GetInformerFactoryByExecutor(caaf.kubernetesClient, fv1.ExecutorTypeContainer)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -44,6 +44,7 @@ import (
|
|||||||
"github.com/fission/fission/pkg/executor/fscache"
|
"github.com/fission/fission/pkg/executor/fscache"
|
||||||
"github.com/fission/fission/pkg/executor/reaper"
|
"github.com/fission/fission/pkg/executor/reaper"
|
||||||
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
|
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
|
||||||
|
finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1"
|
||||||
"github.com/fission/fission/pkg/throttler"
|
"github.com/fission/fission/pkg/throttler"
|
||||||
"github.com/fission/fission/pkg/utils"
|
"github.com/fission/fission/pkg/utils"
|
||||||
"github.com/fission/fission/pkg/utils/maps"
|
"github.com/fission/fission/pkg/utils/maps"
|
||||||
@@ -69,9 +70,6 @@ type (
|
|||||||
|
|
||||||
throttler *throttler.Throttler
|
throttler *throttler.Throttler
|
||||||
|
|
||||||
funcInformer *k8sCache.SharedIndexInformer
|
|
||||||
envInformer *k8sCache.SharedIndexInformer
|
|
||||||
|
|
||||||
serviceInformer k8sCache.SharedIndexInformer
|
serviceInformer k8sCache.SharedIndexInformer
|
||||||
deploymentInformer k8sCache.SharedIndexInformer
|
deploymentInformer k8sCache.SharedIndexInformer
|
||||||
|
|
||||||
@@ -87,8 +85,8 @@ func MakeNewDeploy(
|
|||||||
namespace string,
|
namespace string,
|
||||||
fetcherConfig *fetcherConfig.Config,
|
fetcherConfig *fetcherConfig.Config,
|
||||||
instanceID string,
|
instanceID string,
|
||||||
funcInformer *k8sCache.SharedIndexInformer,
|
funcInformer finformerv1.FunctionInformer,
|
||||||
envInformer *k8sCache.SharedIndexInformer,
|
envInformer finformerv1.EnvironmentInformer,
|
||||||
) (executortype.ExecutorType, error) {
|
) (executortype.ExecutorType, error) {
|
||||||
enableIstio := false
|
enableIstio := false
|
||||||
if len(os.Getenv("ENABLE_ISTIO")) > 0 {
|
if len(os.Getenv("ENABLE_ISTIO")) > 0 {
|
||||||
@@ -115,12 +113,10 @@ func MakeNewDeploy(
|
|||||||
useIstio: enableIstio,
|
useIstio: enableIstio,
|
||||||
|
|
||||||
defaultIdlePodReapTime: 2 * time.Minute,
|
defaultIdlePodReapTime: 2 * time.Minute,
|
||||||
funcInformer: funcInformer,
|
|
||||||
envInformer: envInformer,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
(*nd.funcInformer).AddEventHandler(nd.FunctionEventHandlers())
|
funcInformer.Informer().AddEventHandler(nd.FunctionEventHandlers())
|
||||||
(*nd.envInformer).AddEventHandler(nd.EnvEventHandlers())
|
envInformer.Informer().AddEventHandler(nd.EnvEventHandlers())
|
||||||
|
|
||||||
informerFactory, err := utils.GetInformerFactoryByExecutor(nd.kubernetesClient, fv1.ExecutorTypePoolmgr)
|
informerFactory, err := utils.GetInformerFactoryByExecutor(nd.kubernetesClient, fv1.ExecutorTypePoolmgr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -46,6 +46,7 @@ import (
|
|||||||
"github.com/fission/fission/pkg/executor/fscache"
|
"github.com/fission/fission/pkg/executor/fscache"
|
||||||
"github.com/fission/fission/pkg/executor/reaper"
|
"github.com/fission/fission/pkg/executor/reaper"
|
||||||
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
|
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
|
||||||
|
finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1"
|
||||||
"github.com/fission/fission/pkg/utils"
|
"github.com/fission/fission/pkg/utils"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -76,9 +77,6 @@ type (
|
|||||||
enableIstio bool
|
enableIstio bool
|
||||||
fetcherConfig *fetcherConfig.Config
|
fetcherConfig *fetcherConfig.Config
|
||||||
|
|
||||||
funcInformer *k8sCache.SharedIndexInformer
|
|
||||||
pkgInformer *k8sCache.SharedIndexInformer
|
|
||||||
|
|
||||||
podInformer k8sCache.SharedIndexInformer
|
podInformer k8sCache.SharedIndexInformer
|
||||||
|
|
||||||
defaultIdlePodReapTime time.Duration
|
defaultIdlePodReapTime time.Duration
|
||||||
@@ -104,8 +102,8 @@ func MakeGenericPoolManager(
|
|||||||
functionNamespace string,
|
functionNamespace string,
|
||||||
fetcherConfig *fetcherConfig.Config,
|
fetcherConfig *fetcherConfig.Config,
|
||||||
instanceID string,
|
instanceID string,
|
||||||
funcInformer *k8sCache.SharedIndexInformer,
|
funcInformer finformerv1.FunctionInformer,
|
||||||
pkgInformer *k8sCache.SharedIndexInformer,
|
pkgInformer finformerv1.PackageInformer,
|
||||||
) (executortype.ExecutorType, error) {
|
) (executortype.ExecutorType, error) {
|
||||||
|
|
||||||
gpmLogger := logger.Named("generic_pool_manager")
|
gpmLogger := logger.Named("generic_pool_manager")
|
||||||
@@ -133,14 +131,12 @@ func MakeGenericPoolManager(
|
|||||||
defaultIdlePodReapTime: 2 * time.Minute,
|
defaultIdlePodReapTime: 2 * time.Minute,
|
||||||
fetcherConfig: fetcherConfig,
|
fetcherConfig: fetcherConfig,
|
||||||
enableIstio: enableIstio,
|
enableIstio: enableIstio,
|
||||||
funcInformer: funcInformer,
|
|
||||||
pkgInformer: pkgInformer,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
go gpm.service()
|
go gpm.service()
|
||||||
|
|
||||||
(*gpm.funcInformer).AddEventHandler(FunctionEventHandlers(gpm.logger, gpm.kubernetesClient, gpm.namespace, gpm.enableIstio))
|
funcInformer.Informer().AddEventHandler(FunctionEventHandlers(gpm.logger, gpm.kubernetesClient, gpm.namespace, gpm.enableIstio))
|
||||||
(*gpm.pkgInformer).AddEventHandler(PackageEventHandlers(gpm.logger, gpm.kubernetesClient, gpm.namespace))
|
pkgInformer.Informer().AddEventHandler(PackageEventHandlers(gpm.logger, gpm.kubernetesClient, gpm.namespace))
|
||||||
|
|
||||||
kubeInformerFactory, err := utils.GetInformerFactoryByExecutor(gpm.kubernetesClient, fv1.ExecutorTypePoolmgr)
|
kubeInformerFactory, err := utils.GetInformerFactoryByExecutor(gpm.kubernetesClient, fv1.ExecutorTypePoolmgr)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
Reference in New Issue
Block a user