diff --git a/executor/cleanup.go b/executor/cleanup.go index 558e4d91..6eb2f0b0 100644 --- a/executor/cleanup.go +++ b/executor/cleanup.go @@ -202,6 +202,14 @@ func cleanupDeployments(client *kubernetes.Clientset, namespace string, instance logErr("cleaning up deployment", err) // ignore err } + // Backward compatibility with older label name + pid, pok := dep.ObjectMeta.Labels[fission.POOLMGR_INSTANCEID_LABEL] + if pok && pid != instanceId { + log.Printf("Cleaning up deployment %v", dep.ObjectMeta.Name) + err := client.ExtensionsV1beta1().Deployments(namespace).Delete(dep.ObjectMeta.Name, nil) + logErr("cleaning up deployment", err) + // ignore err + } } return nil } @@ -218,6 +226,13 @@ func cleanupReplicaSets(client *kubernetes.Clientset, namespace string, instance err := client.ExtensionsV1beta1().ReplicaSets(namespace).Delete(rs.ObjectMeta.Name, nil) logErr("cleaning up replicaset", err) } + // Backward compatibility with older label name + pid, pok := rs.ObjectMeta.Labels[fission.POOLMGR_INSTANCEID_LABEL] + if pok && pid != instanceId { + log.Printf("Cleaning up replicaset %v", rs.ObjectMeta.Name) + err := client.ExtensionsV1beta1().ReplicaSets(namespace).Delete(rs.ObjectMeta.Name, nil) + logErr("cleaning up replicaset", err) + } } return nil } @@ -235,6 +250,15 @@ func cleanupPods(client *kubernetes.Clientset, namespace string, instanceId stri logErr("cleaning up pod", err) // ignore err } + // Backward compatibility with older label name + pid, pok := pod.ObjectMeta.Labels[fission.POOLMGR_INSTANCEID_LABEL] + if pok && pid != instanceId { + log.Printf("Cleaning up pod %v", pod.ObjectMeta.Name) + err := client.CoreV1().Pods(namespace).Delete(pod.ObjectMeta.Name, nil) + logErr("cleaning up pod", err) + // ignore err + } + } return nil } diff --git a/executor/poolmgr/gpm.go b/executor/poolmgr/gpm.go index a3f2f7c4..9bbc62d1 100644 --- a/executor/poolmgr/gpm.go +++ b/executor/poolmgr/gpm.go @@ -89,14 +89,14 @@ func (gpm *GenericPoolManager) service() { var err error pool, ok := gpm.pools[crd.CacheKey(&req.env.Metadata)] if !ok { - var poolSize = int32(req.env.Spec.Poolsize) + poolsize := gpm.getEnvPoolsize(req.env) switch req.env.Spec.AllowedFunctionsPerContainer { case fission.AllowedFunctionsPerContainerInfinite: - poolSize = 1 + poolsize = 1 } pool, err = MakeGenericPool( - gpm.fissionClient, gpm.kubernetesClient, req.env, poolSize, + gpm.fissionClient, gpm.kubernetesClient, req.env, poolsize, gpm.namespace, gpm.fsCache, gpm.instanceId) if err != nil { req.responseChannel <- &response{error: err} @@ -110,7 +110,7 @@ func (gpm *GenericPoolManager) service() { latestEnvPoolsize := make(map[string]int) for _, env := range req.envList { latestEnvSet[crd.CacheKey(&env.Metadata)] = true - latestEnvPoolsize[crd.CacheKey(&env.Metadata)] = env.Spec.Poolsize + latestEnvPoolsize[crd.CacheKey(&env.Metadata)] = int(gpm.getEnvPoolsize(&env)) } for key, pool := range gpm.pools { _, ok := latestEnvSet[key] @@ -166,7 +166,7 @@ func (gpm *GenericPoolManager) eagerPoolCreator() { for i := range envs.Items { env := envs.Items[i] // Create pool only if poolsize greater than zero - if env.Spec.Poolsize > 0 { + if gpm.getEnvPoolsize(&env) > 0 { _, err := gpm.GetPool(&envs.Items[i]) if err != nil { log.Printf("eager-create pool failed: %v", err) @@ -178,3 +178,13 @@ func (gpm *GenericPoolManager) eagerPoolCreator() { gpm.CleanupPools(envs.Items) } } + +func (gpm *GenericPoolManager) getEnvPoolsize(env *crd.Environment) int32 { + var poolsize int32 + if env.Spec.Version < 3 { + poolsize = 3 + } else { + poolsize = int32(env.Spec.Poolsize) + } + return poolsize +} diff --git a/types.go b/types.go index c76f10ef..e627d363 100644 --- a/types.go +++ b/types.go @@ -314,6 +314,7 @@ type ( ) const EXECUTOR_INSTANCEID_LABEL string = "executorInstanceId" +const POOLMGR_INSTANCEID_LABEL string = "poolmgrInstanceId" const ( ChecksumTypeSHA256 ChecksumType = "sha256"