K8s informer to work with specific namespaces for executor (#2651)

Consider specific namespaces mentioned by the user in building informers in the executor
- Confimaps
- Secrets
- Deployments
- Services
- Pods
- Replicasets
This commit is contained in:
Shubham Bansal
2022-12-04 21:01:31 +05:30
committed by GitHub
parent 6bf0c4124a
commit 691feaa84f
9 changed files with 180 additions and 135 deletions
+5 -5
View File
@@ -21,8 +21,8 @@ import (
"github.com/pkg/errors"
"go.uber.org/zap"
informerv1 "k8s.io/client-go/informers/core/v1"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/executor/executortype"
@@ -41,18 +41,18 @@ type (
// MakeConfigSecretController makes a controller for configmaps and secrets which changes related functions
func MakeConfigSecretController(ctx context.Context, logger *zap.Logger, fissionClient versioned.Interface,
kubernetesClient kubernetes.Interface, types map[fv1.ExecutorType]executortype.ExecutorType,
configmapInformer map[string]informerv1.ConfigMapInformer,
secretInformer map[string]informerv1.SecretInformer) *ConfigSecretController {
configmapInformer,
secretInformer map[string]cache.SharedIndexInformer) *ConfigSecretController {
logger.Debug("Creating ConfigMap & Secret Controller")
cmsController := &ConfigSecretController{
logger: logger,
fissionClient: fissionClient,
}
for _, informer := range configmapInformer {
informer.Informer().AddEventHandler(ConfigMapEventHandlers(ctx, logger, fissionClient, kubernetesClient, types))
informer.AddEventHandler(ConfigMapEventHandlers(ctx, logger, fissionClient, kubernetesClient, types))
}
for _, informer := range secretInformer {
informer.Informer().AddEventHandler(SecretEventHandlers(ctx, logger, fissionClient, kubernetesClient, types))
informer.AddEventHandler(SecretEventHandlers(ctx, logger, fissionClient, kubernetesClient, types))
}
return cmsController
+25 -33
View File
@@ -29,8 +29,6 @@ import (
"github.com/pkg/errors"
"go.uber.org/zap"
apiv1 "k8s.io/api/core/v1"
k8sInformers "k8s.io/client-go/informers"
k8sInformersv1 "k8s.io/client-go/informers/core/v1"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
@@ -297,49 +295,46 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error {
pkgInformer[ns] = factory.Core().V1().Packages()
}
gpmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypePoolmgr, time.Minute*30)
executorLabel, err := utils.GetInformerLabelByExecutor(fv1.ExecutorTypePoolmgr)
if err != nil {
return err
}
gpmPodInformer := gpmInformerFactory.Core().V1().Pods()
gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets()
gpmInformerFactory := utils.GetInformerFactoryByExecutor(kubernetesClient, executorLabel, time.Minute*30)
gpm, err := poolmgr.MakeGenericPoolManager(ctx,
logger,
fissionClient, kubernetesClient, metricsClient,
fetcherConfig, executorInstanceID,
funcInformer, pkgInformer, envInformer,
gpmPodInformer, gpmRsInformer, podSpecPatch)
gpmInformerFactory, podSpecPatch)
if err != nil {
return errors.Wrap(err, "pool manager creation failed")
}
ndmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeNewdeploy, time.Minute*30)
executorLabel, err = utils.GetInformerLabelByExecutor(fv1.ExecutorTypeNewdeploy)
if err != nil {
return err
}
ndmDeplInformer := ndmInformerFactory.Apps().V1().Deployments()
ndmSvcInformer := ndmInformerFactory.Core().V1().Services()
ndmInformerFactory := utils.GetInformerFactoryByExecutor(kubernetesClient, executorLabel, time.Minute*30)
ndm, err := newdeploy.MakeNewDeploy(ctx,
logger,
fissionClient, kubernetesClient,
fetcherConfig, executorInstanceID,
funcInformer, envInformer,
ndmDeplInformer, ndmSvcInformer, podSpecPatch)
ndmInformerFactory, podSpecPatch)
if err != nil {
return errors.Wrap(err, "new deploy manager creation failed")
}
cnmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeContainer, time.Minute*30)
executorLabel, err = utils.GetInformerLabelByExecutor(fv1.ExecutorTypeContainer)
if err != nil {
return err
}
cnmDeplInformer := cnmInformerFactory.Apps().V1().Deployments()
cnmSvcInformer := cnmInformerFactory.Core().V1().Services()
cnmInformerFactory := utils.GetInformerFactoryByExecutor(kubernetesClient, executorLabel, time.Minute*30)
cnm, err := container.MakeContainer(
ctx, logger,
fissionClient, kubernetesClient,
executorInstanceID, funcInformer,
cnmDeplInformer, cnmSvcInformer)
cnmInformerFactory)
if err != nil {
return errors.Wrap(err, "container manager creation failed")
}
@@ -366,15 +361,8 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error {
// TODO: use context to control the waiting time once kubernetes client supports it.
util.WaitTimeout(wg, 30*time.Second)
configMapInformer := make(map[string]k8sInformersv1.ConfigMapInformer, 0)
secretInformer := make(map[string]k8sInformersv1.SecretInformer, 0)
for _, ns := range utils.DefaultNSResolver().FissionResourceNS {
factory := k8sInformers.NewFilteredSharedInformerFactory(kubernetesClient, time.Minute*30, ns, nil)
configMapInformer[ns] = factory.Core().V1().ConfigMaps()
secretInformer[ns] = factory.Core().V1().Secrets()
}
configMapInformer := utils.GetK8sInformersForNamespaces(kubernetesClient, time.Minute*30, fv1.ConfigMaps)
secretInformer := utils.GetK8sInformersForNamespaces(kubernetesClient, time.Minute*30, fv1.Secrets)
cms := cms.MakeConfigSecretController(ctx, logger, fissionClient, kubernetesClient, executorTypes, configMapInformer, secretInformer)
fissionInformers := make([]k8sCache.SharedIndexInformer, 0)
@@ -388,20 +376,24 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error {
fissionInformers = append(fissionInformers, informer.Informer())
}
for _, informer := range configMapInformer {
fissionInformers = append(fissionInformers, informer.Informer())
fissionInformers = append(fissionInformers, informer)
}
for _, informer := range secretInformer {
fissionInformers = append(fissionInformers, informer.Informer())
fissionInformers = append(fissionInformers, informer)
}
for _, informerFactory := range gpmInformerFactory {
fissionInformers = append(fissionInformers, informerFactory.Core().V1().Pods().Informer())
fissionInformers = append(fissionInformers, informerFactory.Apps().V1().ReplicaSets().Informer())
}
for _, informerFactory := range ndmInformerFactory {
fissionInformers = append(fissionInformers, informerFactory.Apps().V1().Deployments().Informer())
fissionInformers = append(fissionInformers, informerFactory.Core().V1().Services().Informer())
}
for _, informerFactory := range cnmInformerFactory {
fissionInformers = append(fissionInformers, informerFactory.Apps().V1().Deployments().Informer())
fissionInformers = append(fissionInformers, informerFactory.Core().V1().Services().Informer())
}
fissionInformers = append(fissionInformers,
gpmPodInformer.Informer(),
gpmRsInformer.Informer(),
ndmDeplInformer.Informer(),
ndmSvcInformer.Informer(),
cnmDeplInformer.Informer(),
cnmSvcInformer.Informer(),
)
api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes,
fissionInformers...,
)
@@ -35,8 +35,7 @@ import (
"k8s.io/apimachinery/pkg/labels"
k8sTypes "k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/wait"
appsinformers "k8s.io/client-go/informers/apps/v1"
coreinformers "k8s.io/client-go/informers/core/v1"
k8sInformers "k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
appslisters "k8s.io/client-go/listers/apps/v1"
corelisters "k8s.io/client-go/listers/core/v1"
@@ -81,11 +80,11 @@ type (
defaultIdlePodReapTime time.Duration
deplLister appslisters.DeploymentLister
svcLister corelisters.ServiceLister
deplLister map[string]appslisters.DeploymentLister
svcLister map[string]corelisters.ServiceLister
deplListerSynced k8sCache.InformerSynced
svcListerSynced k8sCache.InformerSynced
deplListerSynced map[string]k8sCache.InformerSynced
svcListerSynced map[string]k8sCache.InformerSynced
hpaops *hpautils.HpaOperations
objectReaperIntervalSecond time.Duration
@@ -100,8 +99,7 @@ func MakeContainer(
kubernetesClient kubernetes.Interface,
instanceID string,
funcInformer map[string]finformerv1.FunctionInformer,
deplInformer appsinformers.DeploymentInformer,
svcInformer coreinformers.ServiceInformer,
cnmInformerFactory map[string]k8sInformers.SharedInformerFactory,
) (executortype.ExecutorType, error) {
enableIstio := false
if len(os.Getenv("ENABLE_ISTIO")) > 0 {
@@ -129,14 +127,18 @@ func MakeContainer(
defaultIdlePodReapTime: 1 * time.Minute,
objectReaperIntervalSecond: time.Duration(executorUtils.GetObjectReaperInterval(logger, fv1.ExecutorTypeContainer, 5)) * time.Second,
hpaops: hpautils.NewHpaOperations(logger, kubernetesClient, instanceID),
deplLister: make(map[string]appslisters.DeploymentLister),
deplListerSynced: make(map[string]k8sCache.InformerSynced),
svcLister: make(map[string]corelisters.ServiceLister),
svcListerSynced: make(map[string]k8sCache.InformerSynced),
}
caaf.deplLister = deplInformer.Lister()
caaf.deplListerSynced = deplInformer.Informer().HasSynced
caaf.svcLister = svcInformer.Lister()
caaf.svcListerSynced = svcInformer.Informer().HasSynced
for ns, informerFactory := range cnmInformerFactory {
caaf.deplLister[ns] = informerFactory.Apps().V1().Deployments().Lister()
caaf.deplListerSynced[ns] = informerFactory.Apps().V1().Deployments().Informer().HasSynced
caaf.svcLister[ns] = informerFactory.Core().V1().Services().Lister()
caaf.svcListerSynced[ns] = informerFactory.Core().V1().Services().Informer().HasSynced
}
for _, informer := range funcInformer {
informer.Informer().AddEventHandler(caaf.FuncInformerHandler(ctx))
}
@@ -145,7 +147,15 @@ func MakeContainer(
// Run start the function along with an object reaper.
func (caaf *Container) Run(ctx context.Context) {
if ok := k8sCache.WaitForCacheSync(ctx.Done(), caaf.deplListerSynced, caaf.svcListerSynced); !ok {
waitSynced := make([]k8sCache.InformerSynced, 0)
for _, deplListerSynced := range caaf.deplListerSynced {
waitSynced = append(waitSynced, deplListerSynced)
}
for _, svcListerSynced := range caaf.svcListerSynced {
waitSynced = append(waitSynced, svcListerSynced)
}
if ok := k8sCache.WaitForCacheSync(ctx.Done(), waitSynced...); !ok {
caaf.logger.Fatal("failed to wait for caches to sync")
}
go caaf.idleObjectReaper(ctx)
@@ -214,7 +224,7 @@ func (caaf *Container) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) bool
}
for _, obj := range fsvc.KubernetesObjects {
if strings.ToLower(obj.Kind) == "service" {
_, err := caaf.svcLister.Services(obj.Namespace).Get(obj.Name)
_, err := caaf.svcLister[obj.Namespace].Services(obj.Namespace).Get(obj.Name)
if err != nil {
if !k8sErrs.IsNotFound(err) {
logger.Error("error validating function service", zap.String("function", fsvc.Function.Name), zap.Error(err))
@@ -222,7 +232,7 @@ func (caaf *Container) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) bool
return false
}
} else if strings.ToLower(obj.Kind) == "deployment" {
currentDeploy, err := caaf.deplLister.Deployments(obj.Namespace).Get(obj.Name)
currentDeploy, err := caaf.deplLister[obj.Namespace].Deployments(obj.Namespace).Get(obj.Name)
if err != nil {
if !k8sErrs.IsNotFound(err) {
logger.Error("error validating function deployment", zap.String("function", fsvc.Function.Name), zap.Error(err))
@@ -36,8 +36,7 @@ import (
"k8s.io/apimachinery/pkg/labels"
k8sTypes "k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/wait"
appsinformers "k8s.io/client-go/informers/apps/v1"
coreinformers "k8s.io/client-go/informers/core/v1"
k8sInformers "k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
appslisters "k8s.io/client-go/listers/apps/v1"
corelisters "k8s.io/client-go/listers/core/v1"
@@ -83,11 +82,11 @@ type (
defaultIdlePodReapTime time.Duration
deplLister appslisters.DeploymentLister
svcLister corelisters.ServiceLister
deplLister map[string]appslisters.DeploymentLister
svcLister map[string]corelisters.ServiceLister
deplListerSynced k8sCache.InformerSynced
svcListerSynced k8sCache.InformerSynced
deplListerSynced map[string]k8sCache.InformerSynced
svcListerSynced map[string]k8sCache.InformerSynced
hpaops *hpautils.HpaOperations
@@ -106,8 +105,7 @@ func MakeNewDeploy(
instanceID string,
funcInformer map[string]finformerv1.FunctionInformer,
envInformer map[string]finformerv1.EnvironmentInformer,
deplInformer appsinformers.DeploymentInformer,
svcInformer coreinformers.ServiceInformer,
ndmInformerFactory map[string]k8sInformers.SharedInformerFactory,
podSpecPatch *apiv1.PodSpec,
) (executortype.ExecutorType, error) {
enableIstio := false
@@ -137,15 +135,19 @@ func MakeNewDeploy(
objectReaperIntervalSecond: time.Duration(executorUtils.GetObjectReaperInterval(logger, fv1.ExecutorTypeNewdeploy, 5)) * time.Second,
hpaops: hpautils.NewHpaOperations(logger, kubernetesClient, instanceID),
podSpecPatch: podSpecPatch,
podSpecPatch: podSpecPatch,
deplLister: make(map[string]appslisters.DeploymentLister),
deplListerSynced: make(map[string]k8sCache.InformerSynced),
svcLister: make(map[string]corelisters.ServiceLister),
svcListerSynced: make(map[string]k8sCache.InformerSynced),
}
nd.deplLister = deplInformer.Lister()
nd.deplListerSynced = deplInformer.Informer().HasSynced
nd.svcLister = svcInformer.Lister()
nd.svcListerSynced = svcInformer.Informer().HasSynced
for ns, informerFactory := range ndmInformerFactory {
nd.deplLister[ns] = informerFactory.Apps().V1().Deployments().Lister()
nd.deplListerSynced[ns] = informerFactory.Apps().V1().Deployments().Informer().HasSynced
nd.svcLister[ns] = informerFactory.Core().V1().Services().Lister()
nd.svcListerSynced[ns] = informerFactory.Core().V1().Services().Informer().HasSynced
}
for _, fnInformer := range funcInformer {
fnInformer.Informer().AddEventHandler(nd.FunctionEventHandlers(ctx))
}
@@ -157,7 +159,15 @@ func MakeNewDeploy(
// Run start the function and environment controller along with an object reaper.
func (deploy *NewDeploy) Run(ctx context.Context) {
if ok := k8sCache.WaitForCacheSync(ctx.Done(), deploy.deplListerSynced, deploy.svcListerSynced); !ok {
waitSynced := make([]k8sCache.InformerSynced, 0)
for _, deplListerSynced := range deploy.deplListerSynced {
waitSynced = append(waitSynced, deplListerSynced)
}
for _, svcListerSynced := range deploy.svcListerSynced {
waitSynced = append(waitSynced, svcListerSynced)
}
if ok := k8sCache.WaitForCacheSync(ctx.Done(), waitSynced...); !ok {
deploy.logger.Fatal("failed to wait for caches to sync")
}
go deploy.idleObjectReaper(ctx)
@@ -222,7 +232,7 @@ func (deploy *NewDeploy) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) boo
}
for _, obj := range fsvc.KubernetesObjects {
if strings.ToLower(obj.Kind) == "service" {
_, err := deploy.svcLister.Services(obj.Namespace).Get(obj.Name)
_, err := deploy.svcLister[obj.Namespace].Services(obj.Namespace).Get(obj.Name)
if err != nil {
if !k8sErrs.IsNotFound(err) {
logger.Error("error validating function service", zap.String("function", fsvc.Function.Name), zap.Error(err))
@@ -231,7 +241,7 @@ func (deploy *NewDeploy) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) boo
}
} else if strings.ToLower(obj.Kind) == "deployment" {
currentDeploy, err := deploy.deplLister.Deployments(obj.Namespace).Get(obj.Name)
currentDeploy, err := deploy.deplLister[obj.Namespace].Deployments(obj.Namespace).Get(obj.Name)
if err != nil {
if !k8sErrs.IsNotFound(err) {
logger.Error("error validating function deployment", zap.String("function", fsvc.Function.Name), zap.Error(err))
@@ -55,13 +55,12 @@ func TestRefreshFuncPods(t *testing.T) {
envInformer := map[string]finformerv1.EnvironmentInformer{
metav1.NamespaceAll: informerFactory.Core().V1().Environments(),
}
newDeployInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeNewdeploy, time.Minute*30)
if err != nil {
t.Fatalf("Error creating informer factory: %s", err)
}
deployInformer := newDeployInformerFactory.Apps().V1().Deployments()
svcInformer := newDeployInformerFactory.Core().V1().Services()
executorLabel, err := utils.GetInformerLabelByExecutor(fv1.ExecutorTypeNewdeploy)
if err != nil {
t.Fatalf("Error creating labels for informer: %s", err)
}
ndmInformerFactory := utils.GetInformerFactoryByExecutor(kubernetesClient, executorLabel, time.Minute*30)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
@@ -82,7 +81,7 @@ func TestRefreshFuncPods(t *testing.T) {
}
executor, err := MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, fetcherConfig, "test",
funcInformer, envInformer, deployInformer, svcInformer, podSpecPatch)
funcInformer, envInformer, ndmInformerFactory, podSpecPatch)
if err != nil {
t.Fatalf("new deploy manager creation failed: %s", err)
}
@@ -99,15 +98,27 @@ func TestRefreshFuncPods(t *testing.T) {
go ndm.Run(ctx)
t.Log("New deploy manager started")
runInformers(ctx, []k8sCache.SharedIndexInformer{
informer := []k8sCache.SharedIndexInformer{
envInformer[metav1.NamespaceAll].Informer(),
funcInformer[metav1.NamespaceAll].Informer(),
deployInformer.Informer(),
svcInformer.Informer(),
})
}
for _, informerFactory := range ndmInformerFactory {
informer = append(informer, informerFactory.Apps().V1().Deployments().Informer())
informer = append(informer, informerFactory.Core().V1().Services().Informer())
}
runInformers(ctx, informer)
t.Log("Informers required for new deploy manager started")
if ok := k8sCache.WaitForCacheSync(ctx.Done(), ndm.deplListerSynced, ndm.svcListerSynced); !ok {
waitSynced := make([]k8sCache.InformerSynced, 0)
for _, deplListerSynced := range ndm.deplListerSynced {
waitSynced = append(waitSynced, deplListerSynced)
}
for _, svcListerSynced := range ndm.svcListerSynced {
waitSynced = append(waitSynced, svcListerSynced)
}
if ok := k8sCache.WaitForCacheSync(ctx.Done(), waitSynced...); !ok {
t.Fatal("Timed out waiting for caches to sync")
}
+16 -12
View File
@@ -37,8 +37,7 @@ import (
k8sTypes "k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/wait"
"k8s.io/apimachinery/pkg/watch"
appsinformers "k8s.io/client-go/informers/apps/v1"
coreinformers "k8s.io/client-go/informers/core/v1"
k8sInformers "k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
corelisters "k8s.io/client-go/listers/core/v1"
k8sCache "k8s.io/client-go/tools/cache"
@@ -88,10 +87,10 @@ type (
fetcherConfig *fetcherConfig.Config
// podLister can list/get pods from the shared informer's store
podLister corelisters.PodLister
podLister map[string]corelisters.PodLister
// podListerSynced returns true if the pod store has been synced at least once.
podListerSynced k8sCache.InformerSynced
podListerSynced map[string]k8sCache.InformerSynced
defaultIdlePodReapTime time.Duration
@@ -123,8 +122,7 @@ func MakeGenericPoolManager(ctx context.Context,
funcInformer map[string]finformerv1.FunctionInformer,
pkgInformer map[string]finformerv1.PackageInformer,
envInformer map[string]finformerv1.EnvironmentInformer,
podInformer coreinformers.PodInformer,
rsInformer appsinformers.ReplicaSetInformer,
gpmInformerFactory map[string]k8sInformers.SharedInformerFactory,
podSpecPatch *apiv1.PodSpec,
) (executortype.ExecutorType, error) {
@@ -140,7 +138,7 @@ func MakeGenericPoolManager(ctx context.Context,
}
poolPodC := NewPoolPodController(ctx, gpmLogger, kubernetesClient,
enableIstio, funcInformer, pkgInformer, envInformer, rsInformer, podInformer)
enableIstio, funcInformer, pkgInformer, envInformer, gpmInformerFactory)
gpm := &GenericPoolManager{
logger: gpmLogger,
@@ -159,9 +157,13 @@ func MakeGenericPoolManager(ctx context.Context,
poolPodC: poolPodC,
podSpecPatch: podSpecPatch,
objectReaperIntervalSecond: time.Duration(executorUtils.GetObjectReaperInterval(logger, fv1.ExecutorTypePoolmgr, 5)) * time.Second,
podLister: make(map[string]corelisters.PodLister),
podListerSynced: make(map[string]k8sCache.InformerSynced),
}
for ns, informerFactory := range gpmInformerFactory {
gpm.podLister[ns] = informerFactory.Core().V1().Pods().Lister()
gpm.podListerSynced[ns] = informerFactory.Core().V1().Pods().Informer().HasSynced
}
gpm.podLister = podInformer.Lister()
gpm.podListerSynced = podInformer.Informer().HasSynced
gpm.logger.Debug("inside MakeGenericPoolManager")
@@ -169,8 +171,10 @@ func MakeGenericPoolManager(ctx context.Context,
}
func (gpm *GenericPoolManager) Run(ctx context.Context) {
if ok := k8sCache.WaitForCacheSync(ctx.Done(), gpm.podListerSynced); !ok {
gpm.logger.Fatal("failed to wait for caches to sync")
for _, podListerSynced := range gpm.podListerSynced {
if ok := k8sCache.WaitForCacheSync(ctx.Done(), podListerSynced); !ok {
gpm.logger.Fatal("failed to wait for caches to sync")
}
}
go gpm.service()
gpm.poolPodC.InjectGpm(gpm)
@@ -246,7 +250,7 @@ func (gpm *GenericPoolManager) IsValid(ctx context.Context, fsvc *fscache.FuncSv
otelUtils.SpanTrackEvent(ctx, "IsValid", fscache.GetAttributesForFuncSvc(fsvc)...)
for _, obj := range fsvc.KubernetesObjects {
if strings.ToLower(obj.Kind) == "pod" {
pod, err := gpm.podLister.Pods(obj.Namespace).Get(obj.Name)
pod, err := gpm.podLister[obj.Namespace].Pods(obj.Namespace).Get(obj.Name)
if err == nil && utils.IsReadyPod(pod) {
// Normally, the address format is http://[pod-ip]:[port], however, if the
// Istio is enabled the address format changes to http://[svc-name]:[port].
@@ -29,8 +29,7 @@ import (
"k8s.io/apimachinery/pkg/labels"
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
"k8s.io/apimachinery/pkg/util/wait"
appsinformers "k8s.io/client-go/informers/apps/v1"
coreinformers "k8s.io/client-go/informers/core/v1"
k8sInformers "k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
corelisters "k8s.io/client-go/listers/core/v1"
k8sCache "k8s.io/client-go/tools/cache"
@@ -54,10 +53,10 @@ type (
envListerSynced map[string]k8sCache.InformerSynced
// podLister can list/get pods from the shared informer's store
podLister corelisters.PodLister
podLister map[string]corelisters.PodLister
// podListerSynced returns true if the pod store has been synced at least once.
podListerSynced k8sCache.InformerSynced
podListerSynced map[string]k8sCache.InformerSynced
envCreateUpdateQueue workqueue.RateLimitingInterface
envDeleteQueue workqueue.RateLimitingInterface
@@ -74,8 +73,7 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger,
funcInformer map[string]finformerv1.FunctionInformer,
pkgInformer map[string]finformerv1.PackageInformer,
envInformer map[string]finformerv1.EnvironmentInformer,
rsInformer appsinformers.ReplicaSetInformer,
podInformer coreinformers.PodInformer) *PoolPodController {
gpmInformerFactory map[string]k8sInformers.SharedInformerFactory) *PoolPodController {
logger = logger.Named("pool_pod_controller")
p := &PoolPodController{
logger: logger,
@@ -84,6 +82,8 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger,
enableIstio: enableIstio,
envLister: make(map[string]flisterv1.EnvironmentLister, 0),
envListerSynced: make(map[string]k8sCache.InformerSynced, 0),
podLister: make(map[string]corelisters.PodLister),
podListerSynced: make(map[string]k8sCache.InformerSynced),
envCreateUpdateQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvAddUpdateQueue"),
envDeleteQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvDeleteQueue"),
spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"),
@@ -103,14 +103,16 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger,
p.envLister[ns] = informer.Lister()
p.envListerSynced[ns] = informer.Informer().HasSynced
}
rsInformer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
AddFunc: p.handleRSAdd,
UpdateFunc: p.handleRSUpdate,
DeleteFunc: p.handleRSDelete,
})
for ns, informerFactory := range gpmInformerFactory {
informerFactory.Apps().V1().ReplicaSets().Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
AddFunc: p.handleRSAdd,
UpdateFunc: p.handleRSUpdate,
DeleteFunc: p.handleRSDelete,
})
p.podLister[ns] = informerFactory.Core().V1().Pods().Lister()
p.podListerSynced[ns] = informerFactory.Core().V1().Pods().Informer().HasSynced
}
p.podLister = podInformer.Lister()
p.podListerSynced = podInformer.Informer().HasSynced
p.logger.Info("pool pod controller handlers registered")
return p
}
@@ -138,7 +140,7 @@ func (p *PoolPodController) processRS(rs *apps.ReplicaSet) {
return
}
rsLabelMap["managed"] = "false"
specializedPods, err := p.podLister.Pods(rs.Namespace).List(labels.SelectorFromSet(rsLabelMap))
specializedPods, err := p.podLister[rs.Namespace].Pods(rs.Namespace).List(labels.SelectorFromSet(rsLabelMap))
if err != nil {
logger.Error("Failed to list specialized pods", zap.Error(err))
}
@@ -233,7 +235,9 @@ func (p *PoolPodController) Run(ctx context.Context, stopCh <-chan struct{}) {
p.logger.Info("Waiting for informer caches to sync")
waitSynced := make([]k8sCache.InformerSynced, 0)
waitSynced = append(waitSynced, p.podListerSynced)
for _, synced := range p.podListerSynced {
waitSynced = append(waitSynced, synced)
}
for _, synced := range p.envListerSynced {
waitSynced = append(waitSynced, synced)
}
@@ -373,7 +377,8 @@ func (p *PoolPodController) envDeleteQueueProcessFunc(ctx context.Context) bool
p.logger.Debug("env delete request processing")
p.gpm.cleanupPool(ctx, env)
specializePodLables := getSpecializedPodLabels(env)
specializedPods, err := p.podLister.Pods(p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace)).List(labels.SelectorFromSet(specializePodLables))
ns := p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace)
specializedPods, err := p.podLister[ns].Pods(ns).List(labels.SelectorFromSet(specializePodLables))
if err != nil {
p.logger.Error("failed to list specialized pods", zap.Error(err))
p.envDeleteQueue.Forget(obj)
@@ -413,7 +418,7 @@ func (p *PoolPodController) spCleanupPodQueueProcessFunc(ctx context.Context) bo
p.spCleanupPodQueue.Forget(key)
return false
}
pod, err := p.podLister.Pods(namespace).Get(name)
pod, err := p.podLister[namespace].Pods(namespace).Get(name)
if apierrors.IsNotFound(err) {
p.logger.Info("pod not found", zap.String("key", key))
p.spCleanupPodQueue.Forget(key)
@@ -61,19 +61,17 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
metav1.NamespaceAll: informerFactory.Core().V1().Environments(),
}
gpmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypePoolmgr, time.Minute*30)
executorLabel, err := utils.GetInformerLabelByExecutor(fv1.ExecutorTypePoolmgr)
if err != nil {
t.Fatalf("Error creating informer factory: %v", err)
t.Fatalf("Error creating labels for informer: %v", err)
}
gpmPodInformer := gpmInformerFactory.Core().V1().Pods()
gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets()
gpmInformerFactory := utils.GetInformerFactoryByExecutor(kubernetesClient, executorLabel, time.Minute*30)
ppc := NewPoolPodController(ctx, logger, kubernetesClient, false,
funcInformer,
pkgInformer,
envInformer,
gpmRsInformer,
gpmPodInformer)
gpmInformerFactory)
executorInstanceID := strings.ToLower(uniuri.NewLen(8))
metricsClient := metricsclient.NewSimpleClientset()
@@ -86,7 +84,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
fissionClient, kubernetesClient, metricsClient,
fetcherConfig, executorInstanceID,
funcInformer, pkgInformer, envInformer,
gpmPodInformer, gpmRsInformer, nil)
gpmInformerFactory, nil)
if err != nil {
t.Fatalf("Error creating generic pool manager: %v", err)
}
@@ -95,15 +93,18 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
go ppc.Run(ctx, ctx.Done())
podInformer := gpmPodInformer.Informer()
runInformers(ctx, []k8sCache.SharedIndexInformer{
informers := []k8sCache.SharedIndexInformer{
funcInformer[metav1.NamespaceAll].Informer(),
pkgInformer[metav1.NamespaceAll].Informer(),
envInformer[metav1.NamespaceAll].Informer(),
podInformer,
gpmRsInformer.Informer(),
})
}
for _, informerFactory := range gpmInformerFactory {
informers = append(informers, informerFactory.Core().V1().Pods().Informer())
informers = append(informers, informerFactory.Apps().V1().ReplicaSets().Informer())
}
runInformers(ctx, informers)
pod := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{
@@ -124,7 +125,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
found := false
for found == false && time.Since(start) < time.Second*5 {
t.Log("Waiting for pod to be added to pool")
pod, err := ppc.podLister.Pods(pod.Namespace).Get(pod.Name)
pod, err := ppc.podLister[pod.Namespace].Pods(pod.Namespace).Get(pod.Name)
if err == nil {
found = true
t.Logf("Found pod %#v", pod.ObjectMeta)
+19 -7
View File
@@ -69,6 +69,22 @@ func GetK8sInformersForNamespaces(client kubernetes.Interface, defaultSync time.
return informers
}
func GetInformerFactoryByExecutor(client kubernetes.Interface, labels labels.Selector, defaultResync time.Duration) map[string]k8sInformers.SharedInformerFactory {
informerFactory := make(map[string]k8sInformers.SharedInformerFactory)
namespaces := DefaultNSResolver()
for _, ns := range namespaces.FissionNSWithOptions(WithBuilderNs(), WithFunctionNs(), WithDefaultNs()) {
factory := k8sInformers.NewSharedInformerFactoryWithOptions(client, defaultResync,
k8sInformers.WithTweakListOptions(func(options *metav1.ListOptions) {
options.LabelSelector = labels.String()
}),
k8sInformers.WithNamespace(ns))
informerFactory[ns] = factory
}
return informerFactory
}
func GetInformerFactoryByReadyPod(client kubernetes.Interface, namespace string, labelSelector *metav1.LabelSelector) (k8sInformers.SharedInformerFactory, error) {
informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, 0,
k8sInformers.WithNamespace(namespace),
@@ -79,20 +95,16 @@ func GetInformerFactoryByReadyPod(client kubernetes.Interface, namespace string,
return informerFactory, nil
}
func GetInformerFactoryByExecutor(client kubernetes.Interface, executorType fv1.ExecutorType, defaultResync time.Duration) (k8sInformers.SharedInformerFactory, error) {
func GetInformerLabelByExecutor(executorType fv1.ExecutorType) (labels.Selector, error) {
executorLabel, err := labels.NewRequirement(fv1.EXECUTOR_TYPE, selection.DoubleEquals, []string{string(executorType)})
if err != nil {
return nil, err
}
labelSelector := labels.NewSelector()
labelSelector.Add(*executorLabel)
informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, defaultResync,
k8sInformers.WithTweakListOptions(func(options *metav1.ListOptions) {
options.LabelSelector = labelSelector.String()
}))
return informerFactory, nil
}
return labelSelector, nil
}
func SupportedMetricsAPIVersionAvailable(discoveredAPIGroups *metav1.APIGroupList) bool {
var supportedMetricsAPIVersions = []string{
"v1beta1",