Use pool pod controller with env informer (#2161)
Signed-off-by: Sanket Sudake sanketsudake@gmail.com - Use informers and listers in executors - Passing context properly in pool manager executor - Environment updates in the pool manager would not cause updates in the deployment - Environment update minimizing downtime - Use replicaset controller and environment delete triggers to cleanup specialized pods
This commit is contained in:
@@ -85,6 +85,7 @@ rules:
|
||||
resources:
|
||||
- deployments
|
||||
- deployments/scale
|
||||
- replicasets
|
||||
verbs:
|
||||
- '*'
|
||||
- apiGroups:
|
||||
|
||||
@@ -31,3 +31,10 @@ import (
|
||||
func CacheKey(metadata *metav1.ObjectMeta) string {
|
||||
return fmt.Sprintf("%v_%v", metadata.UID, metadata.ResourceVersion)
|
||||
}
|
||||
|
||||
// CacheKeyForUID create a key that uniquely identifies the
|
||||
// of the object. Since resourceVersion changes on every update and
|
||||
// UIDs are unique, we don't use resource version here
|
||||
func CacheKeyUID(metadata *metav1.ObjectMeta) string {
|
||||
return fmt.Sprintf("%v", metadata.UID)
|
||||
}
|
||||
|
||||
+39
-10
@@ -45,6 +45,7 @@ import (
|
||||
"github.com/fission/fission/pkg/executor/util"
|
||||
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
|
||||
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
|
||||
"github.com/fission/fission/pkg/utils"
|
||||
)
|
||||
|
||||
type (
|
||||
@@ -281,32 +282,50 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames
|
||||
pkgInformer := informerFactory.Core().V1().Packages()
|
||||
envInformer := informerFactory.Core().V1().Environments()
|
||||
|
||||
gpmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypePoolmgr, time.Minute*30)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
gpmPodInformer := gpmInformerFactory.Core().V1().Pods()
|
||||
gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets()
|
||||
gpm, err := poolmgr.MakeGenericPoolManager(
|
||||
logger,
|
||||
fissionClient, kubernetesClient, metricsClient,
|
||||
functionNamespace, fetcherConfig, executorInstanceID,
|
||||
funcInformer, pkgInformer,
|
||||
)
|
||||
funcInformer, pkgInformer, envInformer,
|
||||
gpmPodInformer, gpmRsInformer)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "pool manager creation faied")
|
||||
}
|
||||
|
||||
ndmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeNewdeploy, time.Minute*30)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
ndmDeplInformer := ndmInformerFactory.Apps().V1().Deployments()
|
||||
ndmSvcInformer := ndmInformerFactory.Core().V1().Services()
|
||||
ndm, err := newdeploy.MakeNewDeploy(
|
||||
logger,
|
||||
fissionClient, kubernetesClient,
|
||||
functionNamespace, fetcherConfig, executorInstanceID,
|
||||
funcInformer, envInformer,
|
||||
)
|
||||
ndmDeplInformer, ndmSvcInformer)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "new deploy manager creation faied")
|
||||
}
|
||||
|
||||
cnmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypeContainer, time.Minute*30)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
cnmDeplInformer := cnmInformerFactory.Apps().V1().Deployments()
|
||||
cnmSvcInformer := cnmInformerFactory.Core().V1().Services()
|
||||
ctx := context.Background()
|
||||
cnm, err := container.MakeContainer(
|
||||
ctx,
|
||||
logger,
|
||||
ctx, logger,
|
||||
fissionClient, kubernetesClient,
|
||||
functionNamespace, executorInstanceID, funcInformer)
|
||||
functionNamespace, executorInstanceID, funcInformer,
|
||||
cnmDeplInformer, cnmSvcInformer)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "container manager creation faied")
|
||||
}
|
||||
@@ -339,13 +358,23 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames
|
||||
|
||||
cms := cms.MakeConfigSecretController(ctx, logger, fissionClient, kubernetesClient, executorTypes, configmapInformer, secretInformer)
|
||||
|
||||
api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes, []k8sCache.SharedIndexInformer{
|
||||
funcInformer.Informer(), pkgInformer.Informer(), envInformer.Informer(), configmapInformer.Informer(), secretInformer.Informer(),
|
||||
})
|
||||
api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes,
|
||||
[]k8sCache.SharedIndexInformer{
|
||||
funcInformer.Informer(),
|
||||
pkgInformer.Informer(),
|
||||
envInformer.Informer(),
|
||||
configmapInformer.Informer(),
|
||||
secretInformer.Informer(),
|
||||
gpmPodInformer.Informer(),
|
||||
gpmRsInformer.Informer(),
|
||||
ndmDeplInformer.Informer(),
|
||||
ndmSvcInformer.Informer(),
|
||||
cnmDeplInformer.Informer(),
|
||||
cnmSvcInformer.Informer(),
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
go reaper.CleanupRoleBindings(logger, kubernetesClient, fissionClient, functionNamespace, envBuilderNamespace, time.Minute*30)
|
||||
go api.Serve(port, openTracingEnabled)
|
||||
go serveMetric(logger)
|
||||
|
||||
@@ -29,13 +29,16 @@ import (
|
||||
multierror "github.com/hashicorp/go-multierror"
|
||||
"github.com/pkg/errors"
|
||||
"go.uber.org/zap"
|
||||
appsv1 "k8s.io/api/apps/v1"
|
||||
apiv1 "k8s.io/api/core/v1"
|
||||
k8sErrs "k8s.io/apimachinery/pkg/api/errors"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/labels"
|
||||
k8sTypes "k8s.io/apimachinery/pkg/types"
|
||||
appsinformers "k8s.io/client-go/informers/apps/v1"
|
||||
coreinformers "k8s.io/client-go/informers/core/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
appslisters "k8s.io/client-go/listers/apps/v1"
|
||||
corelisters "k8s.io/client-go/listers/core/v1"
|
||||
k8sCache "k8s.io/client-go/tools/cache"
|
||||
|
||||
fv1 "github.com/fission/fission/pkg/apis/core/v1"
|
||||
@@ -69,10 +72,13 @@ type (
|
||||
|
||||
throttler *throttler.Throttler
|
||||
|
||||
serviceInformer k8sCache.SharedIndexInformer
|
||||
deploymentInformer k8sCache.SharedIndexInformer
|
||||
|
||||
defaultIdlePodReapTime time.Duration
|
||||
|
||||
deplLister appslisters.DeploymentLister
|
||||
svcLister corelisters.ServiceLister
|
||||
|
||||
deplListerSynced k8sCache.InformerSynced
|
||||
svcListerSynced k8sCache.InformerSynced
|
||||
}
|
||||
)
|
||||
|
||||
@@ -84,7 +90,10 @@ func MakeContainer(
|
||||
kubernetesClient *kubernetes.Clientset,
|
||||
namespace string,
|
||||
instanceID string,
|
||||
funcInformer finformerv1.FunctionInformer) (executortype.ExecutorType, error) {
|
||||
funcInformer finformerv1.FunctionInformer,
|
||||
deplInformer appsinformers.DeploymentInformer,
|
||||
svcInformer coreinformers.ServiceInformer,
|
||||
) (executortype.ExecutorType, error) {
|
||||
enableIstio := false
|
||||
if len(os.Getenv("ENABLE_ISTIO")) > 0 {
|
||||
istio, err := strconv.ParseBool(os.Getenv("ENABLE_ISTIO"))
|
||||
@@ -110,21 +119,21 @@ func MakeContainer(
|
||||
// Time is set slightly higher than NewDeploy as cold starts are longer for CaaF
|
||||
defaultIdlePodReapTime: 1 * time.Minute,
|
||||
}
|
||||
caaf.deplLister = deplInformer.Lister()
|
||||
caaf.deplListerSynced = deplInformer.Informer().HasSynced
|
||||
|
||||
caaf.svcLister = svcInformer.Lister()
|
||||
caaf.svcListerSynced = svcInformer.Informer().HasSynced
|
||||
|
||||
funcInformer.Informer().AddEventHandler(caaf.FuncInformerHandler(ctx))
|
||||
|
||||
informerFactory, err := utils.GetInformerFactoryByExecutor(caaf.kubernetesClient, fv1.ExecutorTypeContainer)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
caaf.serviceInformer = informerFactory.Core().V1().Services().Informer()
|
||||
caaf.deploymentInformer = informerFactory.Apps().V1().Deployments().Informer()
|
||||
|
||||
return caaf, nil
|
||||
}
|
||||
|
||||
// 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 {
|
||||
caaf.logger.Fatal("failed to wait for caches to sync")
|
||||
}
|
||||
go caaf.idleObjectReaper()
|
||||
}
|
||||
|
||||
@@ -174,40 +183,6 @@ func (caaf *Container) TapService(ctx context.Context, svcHost string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (caaf *Container) getServiceInfo(ctx context.Context, obj apiv1.ObjectReference) (*apiv1.Service, error) {
|
||||
item, exists, err := utils.GetCachedItem(obj, caaf.serviceInformer)
|
||||
|
||||
if err != nil || !exists {
|
||||
caaf.logger.Debug(
|
||||
"Falling back to getting service info from k8s API -- this may cause performance issues for your function.",
|
||||
zap.Bool("exists", exists),
|
||||
zap.Error(err),
|
||||
)
|
||||
service, err := caaf.kubernetesClient.CoreV1().Services(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{})
|
||||
return service, err
|
||||
}
|
||||
|
||||
service := item.(*apiv1.Service)
|
||||
return service, nil
|
||||
}
|
||||
|
||||
func (caaf *Container) getDeploymentInfo(ctx context.Context, obj apiv1.ObjectReference) (*appsv1.Deployment, error) {
|
||||
item, exists, err := utils.GetCachedItem(obj, caaf.deploymentInformer)
|
||||
|
||||
if err != nil || !exists {
|
||||
caaf.logger.Debug(
|
||||
"Falling back to getting deployment info from k8s API -- this may cause performance issues for your function.",
|
||||
zap.Bool("exists", exists),
|
||||
zap.Error(err),
|
||||
)
|
||||
deployment, err := caaf.kubernetesClient.AppsV1().Deployments(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{})
|
||||
return deployment, err
|
||||
}
|
||||
|
||||
deployment := item.(*appsv1.Deployment)
|
||||
return deployment, nil
|
||||
}
|
||||
|
||||
// IsValid does a get on the service address to ensure it's a valid service, then
|
||||
// scale deployment to 1 replica if there are no available replicas for function.
|
||||
// Return true if no error occurs, return false otherwise.
|
||||
@@ -222,16 +197,15 @@ func (caaf *Container) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) bool
|
||||
}
|
||||
for _, obj := range fsvc.KubernetesObjects {
|
||||
if strings.ToLower(obj.Kind) == "service" {
|
||||
_, err := caaf.getServiceInfo(ctx, obj)
|
||||
_, err := caaf.svcLister.Services(obj.Namespace).Get(obj.Name)
|
||||
if err != nil {
|
||||
if !k8sErrs.IsNotFound(err) {
|
||||
caaf.logger.Error("error validating function service", zap.String("function", fsvc.Function.Name), zap.Error(err))
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
} else if strings.ToLower(obj.Kind) == "deployment" {
|
||||
currentDeploy, err := caaf.getDeploymentInfo(ctx, obj)
|
||||
currentDeploy, err := caaf.deplLister.Deployments(obj.Namespace).Get(obj.Name)
|
||||
if err != nil {
|
||||
if !k8sErrs.IsNotFound(err) {
|
||||
caaf.logger.Error("error validating function deployment", zap.String("function", fsvc.Function.Name), zap.Error(err))
|
||||
|
||||
@@ -28,14 +28,17 @@ import (
|
||||
multierror "github.com/hashicorp/go-multierror"
|
||||
"github.com/pkg/errors"
|
||||
"go.uber.org/zap"
|
||||
appsv1 "k8s.io/api/apps/v1"
|
||||
autoscalingv1 "k8s.io/api/autoscaling/v1"
|
||||
apiv1 "k8s.io/api/core/v1"
|
||||
k8sErrs "k8s.io/apimachinery/pkg/api/errors"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/labels"
|
||||
k8sTypes "k8s.io/apimachinery/pkg/types"
|
||||
appsinformers "k8s.io/client-go/informers/apps/v1"
|
||||
coreinformers "k8s.io/client-go/informers/core/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
appslisters "k8s.io/client-go/listers/apps/v1"
|
||||
corelisters "k8s.io/client-go/listers/core/v1"
|
||||
k8sCache "k8s.io/client-go/tools/cache"
|
||||
|
||||
fv1 "github.com/fission/fission/pkg/apis/core/v1"
|
||||
@@ -70,10 +73,13 @@ type (
|
||||
|
||||
throttler *throttler.Throttler
|
||||
|
||||
serviceInformer k8sCache.SharedIndexInformer
|
||||
deploymentInformer k8sCache.SharedIndexInformer
|
||||
|
||||
defaultIdlePodReapTime time.Duration
|
||||
|
||||
deplLister appslisters.DeploymentLister
|
||||
svcLister corelisters.ServiceLister
|
||||
|
||||
deplListerSynced k8sCache.InformerSynced
|
||||
svcListerSynced k8sCache.InformerSynced
|
||||
}
|
||||
)
|
||||
|
||||
@@ -87,6 +93,8 @@ func MakeNewDeploy(
|
||||
instanceID string,
|
||||
funcInformer finformerv1.FunctionInformer,
|
||||
envInformer finformerv1.EnvironmentInformer,
|
||||
deplInformer appsinformers.DeploymentInformer,
|
||||
svcInformer coreinformers.ServiceInformer,
|
||||
) (executortype.ExecutorType, error) {
|
||||
enableIstio := false
|
||||
if len(os.Getenv("ENABLE_ISTIO")) > 0 {
|
||||
@@ -115,22 +123,23 @@ func MakeNewDeploy(
|
||||
defaultIdlePodReapTime: 2 * time.Minute,
|
||||
}
|
||||
|
||||
nd.deplLister = deplInformer.Lister()
|
||||
nd.deplListerSynced = deplInformer.Informer().HasSynced
|
||||
|
||||
nd.svcLister = svcInformer.Lister()
|
||||
nd.svcListerSynced = svcInformer.Informer().HasSynced
|
||||
|
||||
funcInformer.Informer().AddEventHandler(nd.FunctionEventHandlers())
|
||||
envInformer.Informer().AddEventHandler(nd.EnvEventHandlers())
|
||||
|
||||
informerFactory, err := utils.GetInformerFactoryByExecutor(nd.kubernetesClient, fv1.ExecutorTypePoolmgr)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
nd.serviceInformer = informerFactory.Core().V1().Services().Informer()
|
||||
nd.deploymentInformer = informerFactory.Apps().V1().Deployments().Informer()
|
||||
return nd, nil
|
||||
}
|
||||
|
||||
// Run start the function and environment controller along with an object reaper.
|
||||
func (deploy *NewDeploy) Run(ctx context.Context) {
|
||||
go deploy.serviceInformer.Run(ctx.Done())
|
||||
go deploy.deploymentInformer.Run(ctx.Done())
|
||||
if ok := k8sCache.WaitForCacheSync(ctx.Done(), deploy.deplListerSynced, deploy.svcListerSynced); !ok {
|
||||
deploy.logger.Fatal("failed to wait for caches to sync")
|
||||
}
|
||||
go deploy.idleObjectReaper()
|
||||
}
|
||||
|
||||
@@ -177,40 +186,6 @@ func (deploy *NewDeploy) TapService(ctx context.Context, svcHost string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (deploy *NewDeploy) getServiceInfo(ctx context.Context, obj apiv1.ObjectReference) (*apiv1.Service, error) {
|
||||
item, exists, err := utils.GetCachedItem(obj, deploy.serviceInformer)
|
||||
|
||||
if err != nil || !exists {
|
||||
deploy.logger.Debug(
|
||||
"Falling back to getting service info from k8s API -- this may cause performance issues for your function.",
|
||||
zap.Bool("exists", exists),
|
||||
zap.Error(err),
|
||||
)
|
||||
service, err := deploy.kubernetesClient.CoreV1().Services(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{})
|
||||
return service, err
|
||||
}
|
||||
|
||||
service := item.(*apiv1.Service)
|
||||
return service, nil
|
||||
}
|
||||
|
||||
func (deploy *NewDeploy) getDeploymentInfo(ctx context.Context, obj apiv1.ObjectReference) (*appsv1.Deployment, error) {
|
||||
item, exists, err := utils.GetCachedItem(obj, deploy.deploymentInformer)
|
||||
|
||||
if err != nil || !exists {
|
||||
deploy.logger.Debug(
|
||||
"Falling back to getting deployment info from k8s API -- this may cause performance issues for your function.",
|
||||
zap.Bool("exists", exists),
|
||||
zap.Error(err),
|
||||
)
|
||||
deployment, err := deploy.kubernetesClient.AppsV1().Deployments(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{})
|
||||
return deployment, err
|
||||
}
|
||||
|
||||
deployment := item.(*appsv1.Deployment)
|
||||
return deployment, nil
|
||||
}
|
||||
|
||||
// IsValid does a get on the service address to ensure it's a valid service, then
|
||||
// scale deployment to 1 replica if there are no available replicas for function.
|
||||
// Return true if no error occurs, return false otherwise.
|
||||
@@ -225,7 +200,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.getServiceInfo(ctx, obj)
|
||||
_, err := deploy.svcLister.Services(obj.Namespace).Get(obj.Name)
|
||||
if err != nil {
|
||||
if !k8sErrs.IsNotFound(err) {
|
||||
deploy.logger.Error("error validating function service", zap.String("function", fsvc.Function.Name), zap.Error(err))
|
||||
@@ -234,7 +209,7 @@ func (deploy *NewDeploy) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) boo
|
||||
}
|
||||
|
||||
} else if strings.ToLower(obj.Kind) == "deployment" {
|
||||
currentDeploy, err := deploy.getDeploymentInfo(ctx, obj)
|
||||
currentDeploy, err := deploy.deplLister.Deployments(obj.Namespace).Get(obj.Name)
|
||||
if err != nil {
|
||||
if !k8sErrs.IsNotFound(err) {
|
||||
deploy.logger.Error("error validating function deployment", zap.String("function", fsvc.Function.Name), zap.Error(err))
|
||||
|
||||
@@ -27,3 +27,13 @@ func getEnvPoolSize(env *fv1.Environment) int32 {
|
||||
}
|
||||
return poolsize
|
||||
}
|
||||
|
||||
func getSpecializedPodLabels(env *fv1.Environment) map[string]string {
|
||||
specialPodLabels := make(map[string]string)
|
||||
specialPodLabels[fv1.EXECUTOR_TYPE] = string(fv1.ExecutorTypePoolmgr)
|
||||
specialPodLabels[fv1.ENVIRONMENT_NAME] = env.ObjectMeta.Name
|
||||
specialPodLabels[fv1.ENVIRONMENT_NAMESPACE] = env.ObjectMeta.Namespace
|
||||
specialPodLabels[fv1.ENVIRONMENT_UID] = string(env.ObjectMeta.UID)
|
||||
specialPodLabels["managed"] = "false"
|
||||
return specialPodLabels
|
||||
}
|
||||
|
||||
@@ -90,7 +90,7 @@ func MakeGenericPool(
|
||||
fsCache *fscache.FunctionServiceCache,
|
||||
fetcherConfig *fetcherConfig.Config,
|
||||
instanceID string,
|
||||
enableIstio bool) (*GenericPool, error) {
|
||||
enableIstio bool) *GenericPool {
|
||||
|
||||
gpLogger := logger.Named("generic_pool")
|
||||
|
||||
@@ -129,24 +129,25 @@ func MakeGenericPool(
|
||||
|
||||
gp.runtimeImagePullPolicy = utils.GetImagePullPolicy(os.Getenv("RUNTIME_IMAGE_PULL_POLICY"))
|
||||
|
||||
return gp
|
||||
}
|
||||
|
||||
func (gp *GenericPool) setup(ctx context.Context) error {
|
||||
// create fetcher SA in this ns, if not already created
|
||||
err = fetcherConfig.SetupServiceAccount(gp.kubernetesClient, gp.namespace, nil)
|
||||
err := gp.fetcherConfig.SetupServiceAccount(gp.kubernetesClient, gp.namespace, nil)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "error creating fetcher service account in namespace %q", gp.namespace)
|
||||
return errors.Wrapf(err, "error creating fetcher service account in namespace %q", gp.namespace)
|
||||
}
|
||||
|
||||
// Labels for generic deployment/RS/pods.
|
||||
//gp.labelsForPool = gp.getDeployLabels()
|
||||
|
||||
// create the pool
|
||||
err = gp.createPoolDeployment(context.Background(), env)
|
||||
err = gp.createPoolDeployment(ctx, gp.env)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return err
|
||||
}
|
||||
|
||||
go gp.startReadyPodController()
|
||||
go gp.updateCPUUtilizationSvc()
|
||||
return gp, nil
|
||||
return nil
|
||||
}
|
||||
|
||||
func (gp *GenericPool) getEnvironmentPoolLabels(env *fv1.Environment) map[string]string {
|
||||
@@ -575,7 +576,7 @@ func (gp *GenericPool) getPercent(cpuUsage resource.Quantity, percentage float64
|
||||
}
|
||||
|
||||
// destroys the pool -- the deployment, replicaset and pods
|
||||
func (gp *GenericPool) destroy() error {
|
||||
func (gp *GenericPool) destroy(ctx context.Context) error {
|
||||
close(gp.stopReadyPodControllerCh)
|
||||
|
||||
deletePropagation := metav1.DeletePropagationBackground
|
||||
@@ -584,7 +585,7 @@ func (gp *GenericPool) destroy() error {
|
||||
}
|
||||
|
||||
err := gp.kubernetesClient.AppsV1().
|
||||
Deployments(gp.namespace).Delete(context.TODO(), gp.deployment.ObjectMeta.Name, delOpt)
|
||||
Deployments(gp.namespace).Delete(ctx, gp.deployment.ObjectMeta.Name, delOpt)
|
||||
if err != nil {
|
||||
gp.logger.Error("error destroying deployment",
|
||||
zap.Error(err),
|
||||
|
||||
@@ -203,3 +203,41 @@ func (gp *GenericPool) createPoolDeployment(ctx context.Context, env *fv1.Enviro
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (gp *GenericPool) updatePoolDeployment(ctx context.Context, env *fv1.Environment) error {
|
||||
logger := gp.logger.With(zap.String("env", env.Name), zap.String("namespace", env.Namespace))
|
||||
if gp.env.ObjectMeta.ResourceVersion == env.ObjectMeta.ResourceVersion {
|
||||
logger.Debug("env resource version matching with pool env")
|
||||
return nil
|
||||
}
|
||||
newDeployment := gp.deployment.DeepCopy()
|
||||
spec, err := gp.genDeploymentSpec(env)
|
||||
if err != nil {
|
||||
logger.Error("error generating deployment spec", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
newDeployment.Spec = *spec
|
||||
deployMeta := gp.genDeploymentMeta(env)
|
||||
deployMeta.Name = gp.deployment.Name
|
||||
newDeployment.ObjectMeta = deployMeta
|
||||
|
||||
poolsize := getEnvPoolSize(env)
|
||||
switch env.Spec.AllowedFunctionsPerContainer {
|
||||
case fv1.AllowedFunctionsPerContainerInfinite:
|
||||
poolsize = 1
|
||||
}
|
||||
newDeployment.Spec.Replicas = &poolsize
|
||||
|
||||
depl, err := gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Update(ctx, newDeployment, metav1.UpdateOptions{})
|
||||
if err != nil {
|
||||
logger.Error("error updating deployment in kubernetes", zap.Error(err), zap.String("deployment", depl.Name))
|
||||
return err
|
||||
}
|
||||
// possible concurrency issue here as
|
||||
// gp.env and gp.deployment referenced at few places
|
||||
// we can move update pool to gpm.service if required
|
||||
gp.env = env
|
||||
gp.deployment = depl
|
||||
logger.Info("Updated deployment for pool", zap.String("deployment", depl.Name))
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -35,7 +35,10 @@ import (
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
k8sTypes "k8s.io/apimachinery/pkg/types"
|
||||
"k8s.io/apimachinery/pkg/watch"
|
||||
appsinformers "k8s.io/client-go/informers/apps/v1"
|
||||
coreinformers "k8s.io/client-go/informers/core/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
corelisters "k8s.io/client-go/listers/core/v1"
|
||||
k8sCache "k8s.io/client-go/tools/cache"
|
||||
metricsclient "k8s.io/metrics/pkg/client/clientset/versioned"
|
||||
|
||||
@@ -56,7 +59,7 @@ type requestType int
|
||||
|
||||
const (
|
||||
GET_POOL requestType = iota
|
||||
CLEANUP_POOLS
|
||||
CLEANUP_POOL
|
||||
)
|
||||
|
||||
type (
|
||||
@@ -77,14 +80,20 @@ type (
|
||||
enableIstio bool
|
||||
fetcherConfig *fetcherConfig.Config
|
||||
|
||||
podInformer k8sCache.SharedIndexInformer
|
||||
// podLister can list/get pods from the shared informer's store
|
||||
podLister corelisters.PodLister
|
||||
|
||||
// podListerSynced returns true if the pod store has been synced at least once.
|
||||
podListerSynced k8sCache.InformerSynced
|
||||
|
||||
defaultIdlePodReapTime time.Duration
|
||||
|
||||
poolPodC *PoolPodController
|
||||
}
|
||||
request struct {
|
||||
requestType
|
||||
ctx context.Context
|
||||
env *fv1.Environment
|
||||
envList []fv1.Environment
|
||||
responseChannel chan *response
|
||||
}
|
||||
response struct {
|
||||
@@ -104,6 +113,9 @@ func MakeGenericPoolManager(
|
||||
instanceID string,
|
||||
funcInformer finformerv1.FunctionInformer,
|
||||
pkgInformer finformerv1.PackageInformer,
|
||||
envInformer finformerv1.EnvironmentInformer,
|
||||
podInformer coreinformers.PodInformer,
|
||||
rsInformer appsinformers.ReplicaSetInformer,
|
||||
) (executortype.ExecutorType, error) {
|
||||
|
||||
gpmLogger := logger.Named("generic_pool_manager")
|
||||
@@ -117,6 +129,9 @@ func MakeGenericPoolManager(
|
||||
enableIstio = istio
|
||||
}
|
||||
|
||||
poolPodC := NewPoolPodController(gpmLogger, kubernetesClient, functionNamespace,
|
||||
enableIstio, funcInformer, pkgInformer, envInformer, rsInformer)
|
||||
|
||||
gpm := &GenericPoolManager{
|
||||
logger: gpmLogger,
|
||||
pools: make(map[string]*GenericPool),
|
||||
@@ -131,29 +146,24 @@ func MakeGenericPoolManager(
|
||||
defaultIdlePodReapTime: 2 * time.Minute,
|
||||
fetcherConfig: fetcherConfig,
|
||||
enableIstio: enableIstio,
|
||||
poolPodC: poolPodC,
|
||||
}
|
||||
gpm.podLister = podInformer.Lister()
|
||||
gpm.podListerSynced = podInformer.Informer().HasSynced
|
||||
|
||||
go gpm.service()
|
||||
|
||||
funcInformer.Informer().AddEventHandler(FunctionEventHandlers(gpm.logger, gpm.kubernetesClient, gpm.namespace, gpm.enableIstio))
|
||||
pkgInformer.Informer().AddEventHandler(PackageEventHandlers(gpm.logger, gpm.kubernetesClient, gpm.namespace))
|
||||
|
||||
kubeInformerFactory, err := utils.GetInformerFactoryByExecutor(gpm.kubernetesClient, fv1.ExecutorTypePoolmgr)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
gpm.podInformer = kubeInformerFactory.Core().V1().Pods().Informer()
|
||||
return gpm, nil
|
||||
}
|
||||
|
||||
func (gpm *GenericPoolManager) Run(ctx context.Context) {
|
||||
// eagerPoolCreator must run after CleanupOldExecutorObjects.
|
||||
// Otherwise, the poolmanager may wrongly delete the deployment.
|
||||
go gpm.eagerPoolCreator()
|
||||
go gpm.podInformer.Run(ctx.Done())
|
||||
if ok := k8sCache.WaitForCacheSync(ctx.Done(), gpm.podListerSynced); !ok {
|
||||
gpm.logger.Fatal("failed to wait for caches to sync")
|
||||
}
|
||||
go gpm.service()
|
||||
gpm.poolPodC.InjectGpm(gpm)
|
||||
go gpm.WebsocketStartEventChecker(gpm.kubernetesClient)
|
||||
go gpm.NoActiveConnectionEventChecker(gpm.kubernetesClient)
|
||||
go gpm.idleObjectReaper()
|
||||
go gpm.poolPodC.Run(ctx.Done())
|
||||
}
|
||||
|
||||
func (gpm *GenericPoolManager) GetTypeName(ctx context.Context) fv1.ExecutorType {
|
||||
@@ -168,7 +178,7 @@ func (gpm *GenericPoolManager) GetFuncSvc(ctx context.Context, fn *fv1.Function)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
pool, created, err := gpm.getPool(env)
|
||||
pool, created, err := gpm.getPool(ctx, env)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -207,30 +217,12 @@ func (gpm *GenericPoolManager) TapService(ctx context.Context, svcHost string) e
|
||||
return nil
|
||||
}
|
||||
|
||||
func (gpm *GenericPoolManager) getPodInfo(ctx context.Context, obj apiv1.ObjectReference) (*apiv1.Pod, error) {
|
||||
store := gpm.podInformer.GetStore()
|
||||
|
||||
item, exists, err := store.Get(obj)
|
||||
if err != nil || !exists {
|
||||
item, exists, err = store.GetByKey(fmt.Sprintf("%s/%s", obj.Namespace, obj.Name))
|
||||
}
|
||||
|
||||
if err != nil || !exists {
|
||||
gpm.logger.Debug("Falling back to getting pod info from k8s API -- this may cause performance issues for your function.")
|
||||
pod, err := gpm.kubernetesClient.CoreV1().Pods(obj.Namespace).Get(ctx, obj.Name, metav1.GetOptions{})
|
||||
return pod, err
|
||||
}
|
||||
|
||||
pod := item.(*apiv1.Pod)
|
||||
return pod, nil
|
||||
}
|
||||
|
||||
// IsValid checks if pod is not deleted and that it has the address passed as the argument. Also checks that all the
|
||||
// containers in it are reporting a ready status for the healthCheck.
|
||||
func (gpm *GenericPoolManager) IsValid(ctx context.Context, fsvc *fscache.FuncSvc) bool {
|
||||
for _, obj := range fsvc.KubernetesObjects {
|
||||
if strings.ToLower(obj.Kind) == "pod" {
|
||||
pod, err := gpm.getPodInfo(ctx, obj)
|
||||
pod, err := gpm.podLister.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].
|
||||
@@ -258,7 +250,7 @@ func (gpm *GenericPoolManager) RefreshFuncPods(ctx context.Context, logger *zap.
|
||||
return err
|
||||
}
|
||||
|
||||
gp, created, err := gpm.getPool(env)
|
||||
gp, created, err := gpm.getPool(ctx, env)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -314,7 +306,7 @@ func (gpm *GenericPoolManager) AdoptExistingResources(ctx context.Context) {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
_, created, err := gpm.getPool(&env)
|
||||
_, created, err := gpm.getPool(ctx, &env)
|
||||
if err != nil {
|
||||
gpm.logger.Error("adopt pool failed", zap.Error(err))
|
||||
}
|
||||
@@ -461,7 +453,7 @@ func (gpm *GenericPoolManager) service() {
|
||||
// just because they are missing in the cache, we end up creating another duplicate pool.
|
||||
var err error
|
||||
created := false
|
||||
pool, ok := gpm.pools[crd.CacheKey(&req.env.ObjectMeta)]
|
||||
pool, ok := gpm.pools[crd.CacheKeyUID(&req.env.ObjectMeta)]
|
||||
if !ok {
|
||||
// To support backward compatibility, if envs are created in default ns, we go ahead
|
||||
// and create pools in fission-function ns as earlier.
|
||||
@@ -469,43 +461,47 @@ func (gpm *GenericPoolManager) service() {
|
||||
if req.env.ObjectMeta.Namespace != metav1.NamespaceDefault {
|
||||
ns = req.env.ObjectMeta.Namespace
|
||||
}
|
||||
|
||||
pool, err = MakeGenericPool(gpm.logger,
|
||||
gpm.fissionClient, gpm.kubernetesClient, gpm.metricsClient, req.env, ns,
|
||||
gpm.namespace, gpm.fsCache, gpm.fetcherConfig, gpm.instanceID, gpm.enableIstio)
|
||||
pool = MakeGenericPool(gpm.logger, gpm.fissionClient, gpm.kubernetesClient,
|
||||
gpm.metricsClient, req.env, ns, gpm.namespace, gpm.fsCache,
|
||||
gpm.fetcherConfig, gpm.instanceID, gpm.enableIstio)
|
||||
err = pool.setup(req.ctx)
|
||||
if err != nil {
|
||||
req.responseChannel <- &response{error: err}
|
||||
continue
|
||||
}
|
||||
gpm.pools[crd.CacheKey(&req.env.ObjectMeta)] = pool
|
||||
gpm.pools[crd.CacheKeyUID(&req.env.ObjectMeta)] = pool
|
||||
created = true
|
||||
}
|
||||
req.responseChannel <- &response{pool: pool, created: created}
|
||||
case CLEANUP_POOLS:
|
||||
latestEnvPoolsize := make(map[string]int)
|
||||
for _, env := range req.envList {
|
||||
latestEnvPoolsize[crd.CacheKey(&env.ObjectMeta)] = int(getEnvPoolSize(&env))
|
||||
case CLEANUP_POOL:
|
||||
env := *req.env
|
||||
gpm.logger.Info("destroying pool",
|
||||
zap.String("environment", env.ObjectMeta.Name),
|
||||
zap.String("namespace", env.ObjectMeta.Namespace))
|
||||
|
||||
key := crd.CacheKeyUID(&req.env.ObjectMeta)
|
||||
pool, ok := gpm.pools[key]
|
||||
if !ok {
|
||||
gpm.logger.Error("Could not find pool", zap.String("environment", env.ObjectMeta.Name), zap.String("namespace", env.ObjectMeta.Namespace))
|
||||
return
|
||||
}
|
||||
for key, pool := range gpm.pools {
|
||||
poolsize, ok := latestEnvPoolsize[key]
|
||||
if !ok || poolsize == 0 {
|
||||
// Env no longer exists or pool size changed to zero
|
||||
|
||||
gpm.logger.Info("destroying generic pool", zap.Any("environment", pool.env.ObjectMeta))
|
||||
delete(gpm.pools, key)
|
||||
|
||||
// and delete the pool asynchronously.
|
||||
go pool.destroy() //nolint errcheck
|
||||
}
|
||||
delete(gpm.pools, key)
|
||||
err := pool.destroy(req.ctx)
|
||||
if err != nil {
|
||||
gpm.logger.Error("failed to destroy pool",
|
||||
zap.String("environment", env.ObjectMeta.Name),
|
||||
zap.String("namespace", env.ObjectMeta.Namespace),
|
||||
zap.Error(err))
|
||||
}
|
||||
// no response, caller doesn't wait
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (gpm *GenericPoolManager) getPool(env *fv1.Environment) (*GenericPool, bool, error) {
|
||||
func (gpm *GenericPoolManager) getPool(ctx context.Context, env *fv1.Environment) (*GenericPool, bool, error) {
|
||||
c := make(chan *response)
|
||||
gpm.requestChannel <- &request{
|
||||
ctx: ctx,
|
||||
requestType: GET_POOL,
|
||||
env: env,
|
||||
responseChannel: c,
|
||||
@@ -514,10 +510,11 @@ func (gpm *GenericPoolManager) getPool(env *fv1.Environment) (*GenericPool, bool
|
||||
return resp.pool, resp.created, resp.error
|
||||
}
|
||||
|
||||
func (gpm *GenericPoolManager) cleanupPools(envs []fv1.Environment) {
|
||||
func (gpm *GenericPoolManager) cleanupPool(ctx context.Context, env *fv1.Environment) {
|
||||
gpm.requestChannel <- &request{
|
||||
requestType: CLEANUP_POOLS,
|
||||
envList: envs,
|
||||
ctx: ctx,
|
||||
requestType: CLEANUP_POOL,
|
||||
env: env,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -551,53 +548,6 @@ func (gpm *GenericPoolManager) getFunctionEnv(ctx context.Context, fn *fv1.Funct
|
||||
return env, nil
|
||||
}
|
||||
|
||||
func (gpm *GenericPoolManager) eagerPoolCreator() {
|
||||
pollSleep := 2 * time.Second
|
||||
for {
|
||||
// get list of envs from controller
|
||||
envs, err := gpm.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{})
|
||||
if err != nil {
|
||||
if utils.IsNetworkError(err) {
|
||||
gpm.logger.Error("encountered network error, retrying", zap.Error(err))
|
||||
} else {
|
||||
gpm.logger.Error("failed to get environment list", zap.Error(err))
|
||||
}
|
||||
time.Sleep(5 * time.Second)
|
||||
continue
|
||||
}
|
||||
|
||||
// Create pools for all envs. TODO: we should make this a bit less eager, only
|
||||
// creating pools for envs that are actually used by functions. Also we might want
|
||||
// to keep these eagerly created pools smaller than the ones created when there are
|
||||
// actual function calls.
|
||||
|
||||
wg := &sync.WaitGroup{}
|
||||
|
||||
for i := range envs.Items {
|
||||
env := envs.Items[i]
|
||||
// Create pool only if poolsize greater than zero
|
||||
if getEnvPoolSize(&env) > 0 {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
_, created, err := gpm.getPool(&env)
|
||||
if err != nil {
|
||||
gpm.logger.Error("eager-create pool failed", zap.Error(err))
|
||||
}
|
||||
if created {
|
||||
gpm.logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.namespace))
|
||||
}
|
||||
}()
|
||||
}
|
||||
}
|
||||
|
||||
// Clean up pools whose env was deleted
|
||||
gpm.cleanupPools(envs.Items)
|
||||
wg.Wait()
|
||||
time.Sleep(pollSleep)
|
||||
}
|
||||
}
|
||||
|
||||
// idleObjectReaper reaps objects after certain idle time
|
||||
func (gpm *GenericPoolManager) idleObjectReaper() {
|
||||
ctx := context.Background()
|
||||
|
||||
@@ -0,0 +1,414 @@
|
||||
/*
|
||||
Copyright 2021 The Fission Authors.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
package poolmgr
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
apps "k8s.io/api/apps/v1"
|
||||
v1 "k8s.io/api/core/v1"
|
||||
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"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"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
k8sCache "k8s.io/client-go/tools/cache"
|
||||
"k8s.io/client-go/util/workqueue"
|
||||
|
||||
fv1 "github.com/fission/fission/pkg/apis/core/v1"
|
||||
"github.com/fission/fission/pkg/executor/fscache"
|
||||
finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1"
|
||||
flisterv1 "github.com/fission/fission/pkg/generated/listers/core/v1"
|
||||
)
|
||||
|
||||
type (
|
||||
PoolPodController struct {
|
||||
logger *zap.Logger
|
||||
kubernetesClient *kubernetes.Clientset
|
||||
namespace string
|
||||
enableIstio bool
|
||||
|
||||
envLister flisterv1.EnvironmentLister
|
||||
envListerSynced k8sCache.InformerSynced
|
||||
|
||||
envCreateUpdateQueue workqueue.RateLimitingInterface
|
||||
envDeleteQueue workqueue.RateLimitingInterface
|
||||
|
||||
spCleanupPodQueue workqueue.RateLimitingInterface
|
||||
|
||||
gpm *GenericPoolManager
|
||||
}
|
||||
)
|
||||
|
||||
func NewPoolPodController(logger *zap.Logger,
|
||||
kubernetesClient *kubernetes.Clientset,
|
||||
namespace string,
|
||||
enableIstio bool,
|
||||
funcInformer finformerv1.FunctionInformer,
|
||||
pkgInformer finformerv1.PackageInformer,
|
||||
envInformer finformerv1.EnvironmentInformer,
|
||||
rsInformer appsinformers.ReplicaSetInformer) *PoolPodController {
|
||||
p := &PoolPodController{
|
||||
logger: logger,
|
||||
kubernetesClient: kubernetesClient,
|
||||
namespace: namespace,
|
||||
enableIstio: enableIstio,
|
||||
|
||||
envCreateUpdateQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvAddUpdateQueue"),
|
||||
envDeleteQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvDeleteQueue"),
|
||||
spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"),
|
||||
}
|
||||
funcInformer.Informer().AddEventHandler(FunctionEventHandlers(p.logger, p.kubernetesClient, p.namespace, p.enableIstio))
|
||||
pkgInformer.Informer().AddEventHandler(PackageEventHandlers(p.logger, p.kubernetesClient, p.namespace))
|
||||
envInformer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
|
||||
AddFunc: p.enqueueEnvAdd,
|
||||
UpdateFunc: p.enqueueEnvUpdate,
|
||||
DeleteFunc: p.enqueueEnvDelete,
|
||||
})
|
||||
rsInformer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
|
||||
AddFunc: p.handleRSAdd,
|
||||
UpdateFunc: p.handleRSUpdate,
|
||||
DeleteFunc: p.handleRSDelete,
|
||||
})
|
||||
|
||||
p.envLister = envInformer.Lister()
|
||||
p.envListerSynced = envInformer.Informer().HasSynced
|
||||
p.logger.Info("pool pod controller handlers registered")
|
||||
return p
|
||||
}
|
||||
|
||||
func (p *PoolPodController) InjectGpm(gpm *GenericPoolManager) {
|
||||
p.gpm = gpm
|
||||
}
|
||||
|
||||
func IsPodActive(p *v1.Pod) bool {
|
||||
return v1.PodSucceeded != p.Status.Phase &&
|
||||
v1.PodFailed != p.Status.Phase &&
|
||||
p.DeletionTimestamp == nil
|
||||
}
|
||||
|
||||
func (p *PoolPodController) processRS(rs *apps.ReplicaSet) {
|
||||
if *(rs.Spec.Replicas) != 0 {
|
||||
return
|
||||
}
|
||||
logger := p.logger.With(zap.String("rs", rs.Name), zap.String("namespace", rs.Namespace))
|
||||
logger.Debug("replica set has zero replica count")
|
||||
// List all specialized pods and schedule for cleanup
|
||||
rsLabelMap, err := metav1.LabelSelectorAsMap(rs.Spec.Selector)
|
||||
if err != nil {
|
||||
p.logger.Error("Failed to parse label selector", zap.Error(err))
|
||||
return
|
||||
}
|
||||
rsLabelMap["managed"] = "false"
|
||||
specializedPods, err := p.gpm.podLister.Pods(rs.Namespace).List(labels.SelectorFromSet(rsLabelMap))
|
||||
if err != nil {
|
||||
logger.Error("Failed to list specialized pods", zap.Error(err))
|
||||
}
|
||||
if len(specializedPods) == 0 {
|
||||
return
|
||||
}
|
||||
logger.Info("specialized pods identified for cleanup with RS", zap.Int("numPods", len(specializedPods)))
|
||||
for _, pod := range specializedPods {
|
||||
if !IsPodActive(pod) {
|
||||
continue
|
||||
}
|
||||
key, err := k8sCache.MetaNamespaceKeyFunc(pod)
|
||||
if err != nil {
|
||||
logger.Error("Failed to get key for pod", zap.Error(err))
|
||||
continue
|
||||
}
|
||||
p.spCleanupPodQueue.Add(key)
|
||||
}
|
||||
}
|
||||
|
||||
func (p *PoolPodController) handleRSAdd(obj interface{}) {
|
||||
rs, ok := obj.(*apps.ReplicaSet)
|
||||
if !ok {
|
||||
p.logger.Error("unexpected type when adding rs to pool pod controller", zap.Any("obj", obj))
|
||||
return
|
||||
}
|
||||
p.processRS(rs)
|
||||
}
|
||||
|
||||
func (p *PoolPodController) handleRSUpdate(oldObj interface{}, newObj interface{}) {
|
||||
rs, ok := newObj.(*apps.ReplicaSet)
|
||||
if !ok {
|
||||
p.logger.Error("unexpected type when updating rs to pool pod controller", zap.Any("obj", newObj))
|
||||
return
|
||||
}
|
||||
p.processRS(rs)
|
||||
}
|
||||
|
||||
func (p *PoolPodController) handleRSDelete(obj interface{}) {
|
||||
rs, ok := obj.(*apps.ReplicaSet)
|
||||
if !ok {
|
||||
tombstone, ok := obj.(k8sCache.DeletedFinalStateUnknown)
|
||||
if !ok {
|
||||
p.logger.Error("couldnt get object from tombstone", zap.Any("obj", obj))
|
||||
return
|
||||
}
|
||||
rs, ok = tombstone.Obj.(*apps.ReplicaSet)
|
||||
if !ok {
|
||||
p.logger.Error("tombstone contained object that is not a replicaset", zap.Any("obj", obj))
|
||||
return
|
||||
}
|
||||
}
|
||||
p.processRS(rs)
|
||||
}
|
||||
|
||||
func (p *PoolPodController) enqueueEnvAdd(obj interface{}) {
|
||||
key, err := k8sCache.MetaNamespaceKeyFunc(obj)
|
||||
if err != nil {
|
||||
p.logger.Error("error retrieving key from object in poolPodController", zap.Any("obj", obj))
|
||||
return
|
||||
}
|
||||
p.logger.Debug("enqueue env add", zap.String("key", key))
|
||||
p.envCreateUpdateQueue.Add(key)
|
||||
}
|
||||
|
||||
func (p *PoolPodController) enqueueEnvUpdate(oldObj, newObj interface{}) {
|
||||
key, err := k8sCache.MetaNamespaceKeyFunc(newObj)
|
||||
if err != nil {
|
||||
p.logger.Error("error retrieving key from object in poolPodController", zap.Any("obj", key))
|
||||
return
|
||||
}
|
||||
p.logger.Debug("enqueue env update", zap.String("key", key))
|
||||
p.envCreateUpdateQueue.Add(key)
|
||||
}
|
||||
|
||||
func (p *PoolPodController) enqueueEnvDelete(obj interface{}) {
|
||||
env, ok := obj.(*fv1.Environment)
|
||||
if !ok {
|
||||
p.logger.Error("unexpected type when deleting env to pool pod controller", zap.Any("obj", obj))
|
||||
return
|
||||
}
|
||||
p.logger.Debug("enqueue env delete", zap.Any("env", env))
|
||||
p.envDeleteQueue.Add(env)
|
||||
}
|
||||
|
||||
func (p *PoolPodController) Run(stopCh <-chan struct{}) {
|
||||
defer utilruntime.HandleCrash()
|
||||
defer p.envCreateUpdateQueue.ShutDown()
|
||||
defer p.envDeleteQueue.ShutDown()
|
||||
defer p.spCleanupPodQueue.ShutDown()
|
||||
|
||||
// Wait for the caches to be synced before starting workers
|
||||
p.logger.Info("Waiting for informer caches to sync")
|
||||
if ok := k8sCache.WaitForCacheSync(stopCh, p.envListerSynced); !ok {
|
||||
p.logger.Fatal("failed to wait for caches to sync")
|
||||
}
|
||||
for i := 0; i < 4; i++ {
|
||||
go wait.Until(p.workerRun("envCreateUpdate", p.envCreateUpdateQueueProcessFunc), time.Second, stopCh)
|
||||
}
|
||||
go wait.Until(p.workerRun("envDeleteQueue", p.envDeleteQueueProcessFunc), time.Second, stopCh)
|
||||
go wait.Until(p.workerRun("spCleanupPodQueue", p.spCleanupPodQueueProcessFunc), time.Second, stopCh)
|
||||
p.logger.Info("Started workers for poolPodController")
|
||||
<-stopCh
|
||||
p.logger.Info("Shutting down workers for poolPodController")
|
||||
}
|
||||
|
||||
func (p *PoolPodController) workerRun(name string, processFunc func() bool) func() {
|
||||
return func() {
|
||||
p.logger.Debug("Starting worker with func", zap.String("name", name))
|
||||
for {
|
||||
if quit := processFunc(); quit {
|
||||
p.logger.Info("Shutting down worker", zap.String("name", name))
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (p *PoolPodController) envCreateUpdateQueueProcessFunc() bool {
|
||||
maxRetries := 3
|
||||
handleEnv := func(ctx context.Context, env *fv1.Environment) error {
|
||||
log := p.logger.With(zap.String("env", env.ObjectMeta.Name), zap.String("namespace", env.ObjectMeta.Namespace))
|
||||
log.Debug("env reconsile request processing")
|
||||
pool, created, err := p.gpm.getPool(ctx, env)
|
||||
if err != nil {
|
||||
log.Error("error getting pool", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
if created {
|
||||
log.Info("created pool for the environment")
|
||||
return nil
|
||||
}
|
||||
poolsize := getEnvPoolSize(env)
|
||||
if poolsize == 0 {
|
||||
log.Info("pool size is zero")
|
||||
p.gpm.cleanupPool(ctx, env)
|
||||
return nil
|
||||
}
|
||||
err = pool.updatePoolDeployment(ctx, env)
|
||||
if err != nil {
|
||||
log.Error("error updating pool", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
// If any specialized pods are running, those would be
|
||||
// deleted by replicaSet controller.
|
||||
return nil
|
||||
}
|
||||
|
||||
obj, quit := p.envCreateUpdateQueue.Get()
|
||||
if quit {
|
||||
return true
|
||||
}
|
||||
key := obj.(string)
|
||||
defer p.envCreateUpdateQueue.Done(key)
|
||||
|
||||
namespace, name, err := k8sCache.SplitMetaNamespaceKey(key)
|
||||
if err != nil {
|
||||
p.logger.Error("error splitting key", zap.Error(err))
|
||||
p.envCreateUpdateQueue.Forget(key)
|
||||
return false
|
||||
}
|
||||
env, err := p.envLister.Environments(namespace).Get(name)
|
||||
if apierrors.IsNotFound(err) {
|
||||
p.logger.Info("env not found", zap.String("key", key))
|
||||
p.envCreateUpdateQueue.Forget(key)
|
||||
return false
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
if p.envCreateUpdateQueue.NumRequeues(key) < maxRetries {
|
||||
p.envCreateUpdateQueue.AddRateLimited(key)
|
||||
p.logger.Error("error getting env, retrying", zap.Error(err))
|
||||
} else {
|
||||
p.envCreateUpdateQueue.Forget(key)
|
||||
p.logger.Error("error getting env, retrying, max retries reached", zap.Error(err))
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
err = handleEnv(ctx, env)
|
||||
if err != nil {
|
||||
if p.envCreateUpdateQueue.NumRequeues(key) < maxRetries {
|
||||
p.envCreateUpdateQueue.AddRateLimited(key)
|
||||
p.logger.Error("error handling env from envInformer, retrying", zap.String("key", key), zap.Error(err))
|
||||
} else {
|
||||
p.envCreateUpdateQueue.Forget(key)
|
||||
p.logger.Error("error handling env from envInformer, max retries reached", zap.String("key", key), zap.Error(err))
|
||||
}
|
||||
return false
|
||||
}
|
||||
p.envCreateUpdateQueue.Forget(key)
|
||||
return false
|
||||
}
|
||||
|
||||
func (p *PoolPodController) envDeleteQueueProcessFunc() bool {
|
||||
obj, quit := p.envDeleteQueue.Get()
|
||||
if quit {
|
||||
return true
|
||||
}
|
||||
defer p.envDeleteQueue.Done(obj)
|
||||
env, ok := obj.(*fv1.Environment)
|
||||
if !ok {
|
||||
p.logger.Error("unexpected type when deleting env to pool pod controller", zap.Any("obj", obj))
|
||||
p.envDeleteQueue.Forget(obj)
|
||||
return false
|
||||
}
|
||||
ctx := context.Background()
|
||||
p.logger.Debug("env delete request processing")
|
||||
p.gpm.cleanupPool(ctx, env)
|
||||
specializePodLables := getSpecializedPodLabels(env)
|
||||
specializedPods, err := p.gpm.podLister.Pods(p.gpm.namespace).List(labels.SelectorFromSet(specializePodLables))
|
||||
if err != nil {
|
||||
p.logger.Error("failed to list specialized pods", zap.Error(err))
|
||||
p.envDeleteQueue.Forget(obj)
|
||||
return false
|
||||
}
|
||||
if len(specializedPods) == 0 {
|
||||
p.envDeleteQueue.Forget(obj)
|
||||
return false
|
||||
}
|
||||
p.logger.Info("specialized pods identified for cleanup after env delete", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", env.ObjectMeta.Namespace), zap.Int("count", len(specializedPods)))
|
||||
for _, pod := range specializedPods {
|
||||
if !IsPodActive(pod) {
|
||||
continue
|
||||
}
|
||||
key, err := k8sCache.MetaNamespaceKeyFunc(pod)
|
||||
if err != nil {
|
||||
p.logger.Error("Failed to get key for pod", zap.Error(err))
|
||||
continue
|
||||
}
|
||||
p.spCleanupPodQueue.Add(key)
|
||||
}
|
||||
p.envDeleteQueue.Forget(obj)
|
||||
return false
|
||||
}
|
||||
|
||||
func (p *PoolPodController) spCleanupPodQueueProcessFunc() bool {
|
||||
maxRetries := 3
|
||||
obj, quit := p.spCleanupPodQueue.Get()
|
||||
if quit {
|
||||
return true
|
||||
}
|
||||
key := obj.(string)
|
||||
defer p.spCleanupPodQueue.Done(key)
|
||||
namespace, name, err := k8sCache.SplitMetaNamespaceKey(key)
|
||||
if err != nil {
|
||||
p.logger.Error("error splitting key", zap.Error(err))
|
||||
p.spCleanupPodQueue.Forget(key)
|
||||
return false
|
||||
}
|
||||
pod, err := p.gpm.podLister.Pods(namespace).Get(name)
|
||||
if apierrors.IsNotFound(err) {
|
||||
p.logger.Info("pod not found", zap.String("key", key))
|
||||
p.spCleanupPodQueue.Forget(key)
|
||||
return false
|
||||
}
|
||||
if !IsPodActive(pod) {
|
||||
p.logger.Info("pod not active", zap.String("key", key))
|
||||
p.spCleanupPodQueue.Forget(key)
|
||||
return false
|
||||
}
|
||||
if err != nil {
|
||||
if p.spCleanupPodQueue.NumRequeues(key) < maxRetries {
|
||||
p.spCleanupPodQueue.AddRateLimited(key)
|
||||
p.logger.Error("error getting pod, retrying", zap.Error(err))
|
||||
} else {
|
||||
p.spCleanupPodQueue.Forget(key)
|
||||
p.logger.Error("error getting pod, max retries reached", zap.Error(err))
|
||||
}
|
||||
return false
|
||||
}
|
||||
podName := strings.SplitAfter(pod.GetName(), ".")
|
||||
if fsvc, ok := p.gpm.fsCache.PodToFsvc.Load(strings.TrimSuffix(podName[0], ".")); ok {
|
||||
fsvc, ok := fsvc.(*fscache.FuncSvc)
|
||||
if ok {
|
||||
p.gpm.fsCache.DeleteFunctionSvc(fsvc)
|
||||
p.gpm.fsCache.DeleteEntry(fsvc)
|
||||
} else {
|
||||
p.logger.Error("could not covert item from PodToFsvc", zap.String("key", key))
|
||||
}
|
||||
}
|
||||
err = p.kubernetesClient.CoreV1().Pods(p.namespace).Delete(context.TODO(), pod.Name, metav1.DeleteOptions{})
|
||||
if err != nil {
|
||||
p.logger.Error("failed to delete pod", zap.Error(err), zap.String("pod", pod.ObjectMeta.Name), zap.String("pod_namespace", pod.ObjectMeta.Namespace))
|
||||
return false
|
||||
}
|
||||
p.logger.Info("cleaned specialized pod as environment update/deleted",
|
||||
zap.String("pod", pod.ObjectMeta.Name), zap.String("pod_namespace", pod.ObjectMeta.Namespace),
|
||||
zap.String("address", pod.Status.PodIP))
|
||||
p.spCleanupPodQueue.Forget(key)
|
||||
return false
|
||||
}
|
||||
+3
-16
@@ -1,28 +1,26 @@
|
||||
package utils
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
apiv1 "k8s.io/api/core/v1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/labels"
|
||||
"k8s.io/apimachinery/pkg/selection"
|
||||
k8sInformers "k8s.io/client-go/informers"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
k8sCache "k8s.io/client-go/tools/cache"
|
||||
metricsapi "k8s.io/metrics/pkg/apis/metrics"
|
||||
|
||||
v1 "github.com/fission/fission/pkg/apis/core/v1"
|
||||
)
|
||||
|
||||
func GetInformerFactoryByExecutor(client *kubernetes.Clientset, executorType v1.ExecutorType) (k8sInformers.SharedInformerFactory, error) {
|
||||
func GetInformerFactoryByExecutor(client *kubernetes.Clientset, executorType v1.ExecutorType, defaultResync time.Duration) (k8sInformers.SharedInformerFactory, error) {
|
||||
executorLabel, err := labels.NewRequirement(v1.EXECUTOR_TYPE, selection.DoubleEquals, []string{string(executorType)})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
labelSelector := labels.NewSelector()
|
||||
labelSelector.Add(*executorLabel)
|
||||
informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, 0,
|
||||
informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, defaultResync,
|
||||
k8sInformers.WithTweakListOptions(func(options *metav1.ListOptions) {
|
||||
options.LabelSelector = labelSelector.String()
|
||||
}))
|
||||
@@ -47,14 +45,3 @@ func SupportedMetricsAPIVersionAvailable(discoveredAPIGroups *metav1.APIGroupLis
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func GetCachedItem(obj apiv1.ObjectReference, informer k8sCache.SharedIndexInformer) (item interface{}, exists bool, err error) {
|
||||
store := informer.GetStore()
|
||||
|
||||
item, exists, err = store.Get(obj)
|
||||
if err != nil || !exists {
|
||||
item, exists, err = store.GetByKey(fmt.Sprintf("%s/%s", obj.Namespace, obj.Name))
|
||||
}
|
||||
|
||||
return item, exists, err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user