Pods immediately terminate for idletimeout in new deployment and container executer type (#2459)

* added waitgroup to handle multiple goroutine scenario
* changes for cache service
* changes to get function service from cache by UID
This commit is contained in:
Shubham Bansal
2022-06-29 11:49:34 +05:30
committed by GitHub
parent 473acc4e2b
commit c71867169b
4 changed files with 216 additions and 211 deletions
@@ -34,6 +34,7 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/labels"
k8sTypes "k8s.io/apimachinery/pkg/types" k8sTypes "k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/wait"
appsinformers "k8s.io/client-go/informers/apps/v1" appsinformers "k8s.io/client-go/informers/apps/v1"
coreinformers "k8s.io/client-go/informers/core/v1" coreinformers "k8s.io/client-go/informers/core/v1"
"k8s.io/client-go/kubernetes" "k8s.io/client-go/kubernetes"
@@ -141,7 +142,7 @@ func (caaf *Container) Run(ctx context.Context) {
if ok := k8sCache.WaitForCacheSync(ctx.Done(), caaf.deplListerSynced, caaf.svcListerSynced); !ok { if ok := k8sCache.WaitForCacheSync(ctx.Done(), caaf.deplListerSynced, caaf.svcListerSynced); !ok {
caaf.logger.Fatal("failed to wait for caches to sync") caaf.logger.Fatal("failed to wait for caches to sync")
} }
go caaf.idleObjectReaper() go caaf.idleObjectReaper(ctx)
} }
// GetTypeName returns the executor type name. // GetTypeName returns the executor type name.
@@ -168,7 +169,7 @@ func (caaf *Container) GetFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac
// GetFuncSvcFromCache returns a function service from cache; error otherwise. // GetFuncSvcFromCache returns a function service from cache; error otherwise.
func (caaf *Container) GetFuncSvcFromCache(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) { func (caaf *Container) GetFuncSvcFromCache(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) {
otelUtils.SpanTrackEvent(ctx, "GetFuncSvcFromCache", otelUtils.GetAttributesForFunction(fn)...) otelUtils.SpanTrackEvent(ctx, "GetFuncSvcFromCache", otelUtils.GetAttributesForFunction(fn)...)
return caaf.fsCache.GetByFunction(&fn.ObjectMeta) return caaf.fsCache.GetByFunctionUID(fn.UID)
} }
// DeleteFuncSvcFromCache deletes a function service from cache. // DeleteFuncSvcFromCache deletes a function service from cache.
@@ -703,17 +704,16 @@ func (caaf *Container) updateStatus(fn *fv1.Function, err error, message string)
} }
// idleObjectReaper reaps objects after certain idle time // idleObjectReaper reaps objects after certain idle time
func (caaf *Container) idleObjectReaper() { func (caaf *Container) idleObjectReaper(ctx context.Context) {
ctx := context.Background() // calling function doIdleObjectReaper() repeatedly at given interval of time
wait.UntilWithContext(ctx, caaf.doIdleObjectReaper, time.Second*5)
}
pollSleep := 5 * time.Second func (caaf *Container) doIdleObjectReaper(ctx context.Context) {
for { funcSvcs, err := caaf.fsCache.ListOld(time.Second * 5)
time.Sleep(pollSleep)
funcSvcs, err := caaf.fsCache.ListOld(pollSleep)
if err != nil { if err != nil {
caaf.logger.Error("error reaping idle pods", zap.Error(err)) caaf.logger.Error("error reaping idle pods", zap.Error(err))
continue return
} }
for i := range funcSvcs { for i := range funcSvcs {
@@ -771,7 +771,6 @@ func (caaf *Container) idleObjectReaper() {
}() }()
} }
} }
}
func getDeploymentObj(kubeobjs []apiv1.ObjectReference) *apiv1.ObjectReference { func getDeploymentObj(kubeobjs []apiv1.ObjectReference) *apiv1.ObjectReference {
for _, kubeobj := range kubeobjs { for _, kubeobj := range kubeobjs {
@@ -35,6 +35,7 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/labels"
k8sTypes "k8s.io/apimachinery/pkg/types" k8sTypes "k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/wait"
appsinformers "k8s.io/client-go/informers/apps/v1" appsinformers "k8s.io/client-go/informers/apps/v1"
coreinformers "k8s.io/client-go/informers/core/v1" coreinformers "k8s.io/client-go/informers/core/v1"
"k8s.io/client-go/kubernetes" "k8s.io/client-go/kubernetes"
@@ -169,7 +170,7 @@ func (deploy *NewDeploy) GetFuncSvc(ctx context.Context, fn *fv1.Function) (*fsc
// GetFuncSvcFromCache returns a function service from cache; error otherwise. // GetFuncSvcFromCache returns a function service from cache; error otherwise.
func (deploy *NewDeploy) GetFuncSvcFromCache(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) { func (deploy *NewDeploy) GetFuncSvcFromCache(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) {
otelUtils.SpanTrackEvent(ctx, "GetFuncSvcFromCache") otelUtils.SpanTrackEvent(ctx, "GetFuncSvcFromCache")
return deploy.fsCache.GetByFunction(&fn.ObjectMeta) return deploy.fsCache.GetByFunctionUID(fn.UID)
} }
// DeleteFuncSvcFromCache deletes a function service from cache. // DeleteFuncSvcFromCache deletes a function service from cache.
@@ -593,6 +594,7 @@ func (deploy *NewDeploy) updateFunction(ctx context.Context, oldFn *fv1.Function
if oldFn.Spec.Environment != newFn.Spec.Environment || if oldFn.Spec.Environment != newFn.Spec.Environment ||
oldFn.Spec.Package.PackageRef != newFn.Spec.Package.PackageRef || oldFn.Spec.Package.PackageRef != newFn.Spec.Package.PackageRef ||
oldFn.Spec.Package.FunctionName != newFn.Spec.Package.FunctionName { oldFn.Spec.Package.FunctionName != newFn.Spec.Package.FunctionName {
deploy.logger.Debug("deployment changed", zap.String("msg", "deployment changed"))
deployChanged = true deployChanged = true
} }
@@ -766,10 +768,11 @@ func (deploy *NewDeploy) updateStatus(fn *fv1.Function, err error, message strin
// idleObjectReaper reaps objects after certain idle time // idleObjectReaper reaps objects after certain idle time
func (deploy *NewDeploy) idleObjectReaper(ctx context.Context) { func (deploy *NewDeploy) idleObjectReaper(ctx context.Context) {
pollSleep := 5 * time.Second // calling function doIdleObjectReaper() repeatedly at given interval of time
for { wait.UntilWithContext(ctx, deploy.doIdleObjectReaper, time.Second*5)
time.Sleep(pollSleep) }
func (deploy *NewDeploy) doIdleObjectReaper(ctx context.Context) {
envs, err := deploy.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) envs, err := deploy.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(ctx, metav1.ListOptions{})
if err != nil { if err != nil {
deploy.logger.Fatal("failed to get environment list", zap.Error(err)) deploy.logger.Fatal("failed to get environment list", zap.Error(err))
@@ -780,15 +783,14 @@ func (deploy *NewDeploy) idleObjectReaper(ctx context.Context) {
envList[env.ObjectMeta.UID] = struct{}{} envList[env.ObjectMeta.UID] = struct{}{}
} }
funcSvcs, err := deploy.fsCache.ListOld(pollSleep) funcSvcs, err := deploy.fsCache.ListOld(time.Second * 5)
if err != nil { if err != nil {
deploy.logger.Error("error reaping idle pods", zap.Error(err)) deploy.logger.Error("error reaping idle pods", zap.Error(err))
continue return
} }
for i := range funcSvcs { for i := range funcSvcs {
fsvc := funcSvcs[i] fsvc := funcSvcs[i]
if fsvc.Executor != fv1.ExecutorTypeNewdeploy { if fsvc.Executor != fv1.ExecutorTypeNewdeploy {
continue continue
} }
@@ -849,7 +851,6 @@ func (deploy *NewDeploy) idleObjectReaper(ctx context.Context) {
}() }()
} }
} }
}
func getDeploymentObj(kubeobjs []apiv1.ObjectReference) *apiv1.ObjectReference { func getDeploymentObj(kubeobjs []apiv1.ObjectReference) *apiv1.ObjectReference {
for _, kubeobj := range kubeobjs { for _, kubeobj := range kubeobjs {
+11 -12
View File
@@ -35,6 +35,7 @@ import (
"k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/labels"
"k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime"
k8sTypes "k8s.io/apimachinery/pkg/types" k8sTypes "k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/wait"
"k8s.io/apimachinery/pkg/watch" "k8s.io/apimachinery/pkg/watch"
appsinformers "k8s.io/client-go/informers/apps/v1" appsinformers "k8s.io/client-go/informers/apps/v1"
coreinformers "k8s.io/client-go/informers/core/v1" coreinformers "k8s.io/client-go/informers/core/v1"
@@ -169,7 +170,7 @@ func (gpm *GenericPoolManager) Run(ctx context.Context) {
gpm.poolPodC.InjectGpm(gpm) gpm.poolPodC.InjectGpm(gpm)
go gpm.WebsocketStartEventChecker(gpm.kubernetesClient) go gpm.WebsocketStartEventChecker(gpm.kubernetesClient)
go gpm.NoActiveConnectionEventChecker(gpm.kubernetesClient) go gpm.NoActiveConnectionEventChecker(gpm.kubernetesClient)
go gpm.idleObjectReaper() go gpm.idleObjectReaper(ctx)
go gpm.poolPodC.Run(ctx.Done()) go gpm.poolPodC.Run(ctx.Done())
} }
@@ -569,17 +570,16 @@ func (gpm *GenericPoolManager) getFunctionEnv(ctx context.Context, fn *fv1.Funct
} }
// idleObjectReaper reaps objects after certain idle time // idleObjectReaper reaps objects after certain idle time
func (gpm *GenericPoolManager) idleObjectReaper() { func (gpm *GenericPoolManager) idleObjectReaper(ctx context.Context) {
ctx := context.Background() // calling function doIdleObjectReaper() repeatedly at given interval of time
pollSleep := 5 * time.Second wait.UntilWithContext(ctx, gpm.doIdleObjectReaper, time.Second*5)
}
for {
time.Sleep(pollSleep)
func (gpm *GenericPoolManager) doIdleObjectReaper(ctx context.Context) {
envs, err := gpm.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) envs, err := gpm.fissionClient.CoreV1().Environments(metav1.NamespaceAll).List(ctx, metav1.ListOptions{})
if err != nil { if err != nil {
gpm.logger.Error("failed to get environment list", zap.Error(err)) gpm.logger.Error("failed to get environment list", zap.Error(err))
continue return
} }
envList := make(map[k8sTypes.UID]struct{}) envList := make(map[k8sTypes.UID]struct{})
@@ -590,7 +590,7 @@ func (gpm *GenericPoolManager) idleObjectReaper() {
fns, err := gpm.fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) fns, err := gpm.fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(ctx, metav1.ListOptions{})
if err != nil { if err != nil {
gpm.logger.Error("failed to get environment list", zap.Error(err)) gpm.logger.Error("failed to get environment list", zap.Error(err))
continue return
} }
fnList := make(map[k8sTypes.UID]fv1.Function) fnList := make(map[k8sTypes.UID]fv1.Function)
@@ -598,10 +598,10 @@ func (gpm *GenericPoolManager) idleObjectReaper() {
fnList[fn.ObjectMeta.UID] = fns.Items[i] fnList[fn.ObjectMeta.UID] = fns.Items[i]
} }
funcSvcs, err := gpm.fsCache.ListOldForPool(pollSleep) funcSvcs, err := gpm.fsCache.ListOldForPool(time.Second * 5)
if err != nil { if err != nil {
gpm.logger.Error("error reaping idle pods", zap.Error(err)) gpm.logger.Error("error reaping idle pods", zap.Error(err))
continue return
} }
for i := range funcSvcs { for i := range funcSvcs {
@@ -659,7 +659,6 @@ func (gpm *GenericPoolManager) idleObjectReaper() {
}() }()
} }
} }
}
// WebsocketStartEventChecker checks if the pod has emitted a websocket connection start event // WebsocketStartEventChecker checks if the pod has emitted a websocket connection start event
func (gpm *GenericPoolManager) WebsocketStartEventChecker(kubeClient kubernetes.Interface) { func (gpm *GenericPoolManager) WebsocketStartEventChecker(kubeClient kubernetes.Interface) {
+8 -2
View File
@@ -130,10 +130,16 @@ func (fsc *FunctionServiceCache) service() {
resp.error = fsc._touchByAddress(req.address) resp.error = fsc._touchByAddress(req.address)
case LISTOLD: case LISTOLD:
// get svcs idle for > req.age // get svcs idle for > req.age
fscs := fsc.byFunction.Copy() fscs := fsc.byFunctionUID.Copy()
funcObjects := make([]*FuncSvc, 0) funcObjects := make([]*FuncSvc, 0)
for _, funcSvc := range fscs { for _, funcSvc := range fscs {
fsvc := funcSvc.(*FuncSvc) mI := funcSvc.(metav1.ObjectMeta)
fsvcI, err := fsc.byFunction.Get(crd.CacheKey(&mI))
if err != nil {
fsc.logger.Error("error while getting service", zap.Any("error", err))
return
}
fsvc := fsvcI.(*FuncSvc)
if time.Since(fsvc.Atime) > req.age { if time.Since(fsvc.Atime) > req.age {
funcObjects = append(funcObjects, fsvc) funcObjects = append(funcObjects, fsvc)
} }