Fix namespace used in speciallized pod cleanup (#2415)

In pool pod controller we were using pool namespace
rather than pod namespace in cleanup which was causing
issue in few scenarios. Using pod namespace now instead.
Also add unit test for scenario which was failing.
Using kubernetes client interface now across instead of
kubernetes ClientSet for testing.

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
Sanket Sudake
2022-04-21 14:44:09 +05:30
committed by GitHub
parent 24b185db89
commit 3bdbeb6c87
10 changed files with 188 additions and 37 deletions
@@ -41,7 +41,7 @@ func getIstioServiceLabels(fnName string) map[string]string {
// Based on function create/update/delete event, we create role binding
// for the secret/configmap access which is used by fetcher component.
// If istio is enabled, we create a service for the function.
func FunctionEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clientset, fissionfnNamespace string, istioEnabled bool) k8sCache.ResourceEventHandlerFuncs {
func FunctionEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Interface, fissionfnNamespace string, istioEnabled bool) k8sCache.ResourceEventHandlerFuncs {
return k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
ctx := context.Background()
+8 -7
View File
@@ -49,6 +49,7 @@ import (
"github.com/fission/fission/pkg/executor/metrics"
fetcherClient "github.com/fission/fission/pkg/fetcher/client"
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
"github.com/fission/fission/pkg/generated/clientset/versioned"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/maps"
otelUtils "github.com/fission/fission/pkg/utils/otel"
@@ -67,9 +68,9 @@ type (
useSvc bool // create k8s service for specialized pods
useIstio bool
runtimeImagePullPolicy apiv1.PullPolicy // pull policy for generic pool to created env deployment
kubernetesClient *kubernetes.Clientset
metricsClient *metricsclient.Clientset
fissionClient *crd.FissionClient
kubernetesClient kubernetes.Interface
metricsClient metricsclient.Interface
fissionClient versioned.Interface
fetcherConfig *fetcherConfig.Config
stopReadyPodControllerCh chan struct{}
readyPodLister corelisters.PodLister
@@ -85,9 +86,9 @@ type (
// MakeGenericPool returns an instance of GenericPool
func MakeGenericPool(
logger *zap.Logger,
fissionClient *crd.FissionClient,
kubernetesClient *kubernetes.Clientset,
metricsClient *metricsclient.Clientset,
fissionClient versioned.Interface,
kubernetesClient kubernetes.Interface,
metricsClient metricsclient.Interface,
env *fv1.Environment,
namespace string,
functionNamespace string,
@@ -173,7 +174,7 @@ func (gp *GenericPool) getDeployAnnotations(env *fv1.Environment) map[string]str
}
func (gp *GenericPool) checkMetricsApi() bool {
apiGroups, err := gp.metricsClient.DiscoveryClient.ServerGroups()
apiGroups, err := gp.metricsClient.Discovery().ServerGroups()
if err != nil {
gp.logger.Error("failed to discover API groups", zap.Error(err))
return false
+9 -8
View File
@@ -50,6 +50,7 @@ import (
"github.com/fission/fission/pkg/executor/fscache"
"github.com/fission/fission/pkg/executor/reaper"
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
"github.com/fission/fission/pkg/generated/clientset/versioned"
finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1"
"github.com/fission/fission/pkg/utils"
otelUtils "github.com/fission/fission/pkg/utils/otel"
@@ -69,11 +70,11 @@ type (
logger *zap.Logger
pools map[string]*GenericPool
kubernetesClient *kubernetes.Clientset
metricsClient *metricsclient.Clientset
kubernetesClient kubernetes.Interface
metricsClient metricsclient.Interface
namespace string
fissionClient *crd.FissionClient
fissionClient versioned.Interface
functionEnv *cache.Cache
fsCache *fscache.FunctionServiceCache
instanceID string
@@ -107,9 +108,9 @@ type (
func MakeGenericPoolManager(
logger *zap.Logger,
fissionClient *crd.FissionClient,
kubernetesClient *kubernetes.Clientset,
metricsClient *metricsclient.Clientset,
fissionClient versioned.Interface,
kubernetesClient kubernetes.Interface,
metricsClient metricsclient.Interface,
functionNamespace string,
fetcherConfig *fetcherConfig.Config,
instanceID string,
@@ -657,7 +658,7 @@ func (gpm *GenericPoolManager) idleObjectReaper() {
}
// WebsocketStartEventChecker checks if the pod has emitted a websocket connection start event
func (gpm *GenericPoolManager) WebsocketStartEventChecker(kubeClient *kubernetes.Clientset) {
func (gpm *GenericPoolManager) WebsocketStartEventChecker(kubeClient kubernetes.Interface) {
informer := k8sCache.NewSharedInformer(
&k8sCache.ListWatch{
@@ -698,7 +699,7 @@ func (gpm *GenericPoolManager) WebsocketStartEventChecker(kubeClient *kubernetes
}
// NoActiveConnectionEventChecker checks if the pod has emitted an inactive event
func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(kubeClient *kubernetes.Clientset) {
func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(kubeClient kubernetes.Interface) {
informer := k8sCache.NewSharedInformer(
&k8sCache.ListWatch{
@@ -31,7 +31,7 @@ import (
// PackageEventHandlers provides handlers for package events.
// Based on package create/update event, we create role binding
// for the package which is used by fetcher component.
func PackageEventHandlers(logger *zap.Logger, kubernetesClient *kubernetes.Clientset, fissionfnNamespace string) k8sCache.ResourceEventHandlerFuncs {
func PackageEventHandlers(logger *zap.Logger, kubernetesClient kubernetes.Interface, fissionfnNamespace string) k8sCache.ResourceEventHandlerFuncs {
return k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
pkg := obj.(*fv1.Package)
@@ -44,7 +44,7 @@ import (
type (
PoolPodController struct {
logger *zap.Logger
kubernetesClient *kubernetes.Clientset
kubernetesClient kubernetes.Interface
namespace string
enableIstio bool
@@ -67,7 +67,7 @@ type (
)
func NewPoolPodController(logger *zap.Logger,
kubernetesClient *kubernetes.Clientset,
kubernetesClient kubernetes.Interface,
namespace string,
enableIstio bool,
funcInformer finformerv1.FunctionInformer,
@@ -403,24 +403,24 @@ func (p *PoolPodController) spCleanupPodQueueProcessFunc() bool {
}
return false
}
ctx := context.Background()
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 {
ctx := context.Background()
p.gpm.fsCache.DeleteFunctionSvc(ctx, fsvc)
p.gpm.fsCache.DeleteEntry(fsvc)
} else {
p.logger.Error("could not convert item from PodToFsvc", zap.String("key", key))
}
}
err = p.kubernetesClient.CoreV1().Pods(p.namespace).Delete(context.TODO(), pod.Name, metav1.DeleteOptions{})
err = p.kubernetesClient.CoreV1().Pods(namespace).Delete(ctx, 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))
p.logger.Error("failed to delete pod", zap.Error(err), zap.String("pod", name), zap.String("pod_namespace", 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("pod", name), zap.String("pod_namespace", namespace),
zap.String("address", pod.Status.PodIP))
p.spCleanupPodQueue.Forget(key)
return false
@@ -0,0 +1,149 @@
/*
Copyright 2022 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"
"testing"
"time"
"github.com/dchest/uniuri"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes/fake"
k8sCache "k8s.io/client-go/tools/cache"
metricsclient "k8s.io/metrics/pkg/client/clientset/versioned/fake"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
fClient "github.com/fission/fission/pkg/generated/clientset/versioned/fake"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/loggerfactory"
)
func runInformers(ctx context.Context, informers []k8sCache.SharedIndexInformer) {
// Run all informers
for _, informer := range informers {
go informer.Run(ctx.Done())
}
}
func TestPoolPodControllerPodCleanup(t *testing.T) {
logger := loggerfactory.GetLogger()
kubernetesClient := fake.NewSimpleClientset()
fissionClient := fClient.NewSimpleClientset()
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30)
funcInformer := informerFactory.Core().V1().Functions()
pkgInformer := informerFactory.Core().V1().Packages()
envInformer := informerFactory.Core().V1().Environments()
gpmInformerFactory, err := utils.GetInformerFactoryByExecutor(kubernetesClient, fv1.ExecutorTypePoolmgr, time.Minute*30)
if err != nil {
t.Fatalf("Error creating informer factory: %v", err)
}
gpmPodInformer := gpmInformerFactory.Core().V1().Pods()
gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets()
fnNamespace := "fission-function"
ppc := NewPoolPodController(logger, kubernetesClient, fnNamespace, false,
funcInformer,
pkgInformer,
envInformer,
gpmRsInformer,
gpmPodInformer)
executorInstanceID := strings.ToLower(uniuri.NewLen(8))
metricsClient := metricsclient.NewSimpleClientset()
fetcherConfig, err := fetcherConfig.MakeFetcherConfig("/userfunc")
if err != nil {
t.Fatalf("Error creating fetcher config: %v", err)
}
executor, err := MakeGenericPoolManager(
logger,
fissionClient, kubernetesClient, metricsClient,
fnNamespace, fetcherConfig, executorInstanceID,
funcInformer, pkgInformer, envInformer,
gpmPodInformer, gpmRsInformer)
if err != nil {
t.Fatalf("Error creating generic pool manager: %v", err)
}
gpm := executor.(*GenericPoolManager)
ppc.InjectGpm(gpm)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
go ppc.Run(ctx.Done())
podInformer := gpmPodInformer.Informer()
runInformers(ctx, []k8sCache.SharedIndexInformer{
funcInformer.Informer(),
pkgInformer.Informer(),
envInformer.Informer(),
podInformer,
gpmRsInformer.Informer(),
})
pod := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{
Name: "test-pod",
Namespace: "test-different-namespace",
},
Status: corev1.PodStatus{
Phase: corev1.PodRunning,
},
}
_, err = kubernetesClient.CoreV1().Pods(pod.Namespace).Create(ctx, pod, metav1.CreateOptions{})
if err != nil {
t.Fatalf("Error creating pod: %v", err)
}
// Wait for pod to be added to informer
start := time.Now()
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)
if err == nil {
found = true
t.Logf("Found pod %#v", pod.ObjectMeta)
}
time.Sleep(time.Millisecond * 100)
}
t.Log("Pod added to pool")
// Ask the controller to clean up the pod
key, err := k8sCache.MetaNamespaceKeyFunc(pod)
if err != nil {
t.Fatalf("Error creating key: %v", err)
}
ppc.spCleanupPodQueue.Add(key)
start = time.Now()
for ppc.spCleanupPodQueue.Len() > 0 && time.Since(start) < time.Second*5 {
time.Sleep(time.Millisecond * 100)
t.Log("Waiting for pod cleanup to complete")
}
t.Log("Cleanup pod queue is empty")
// Ensure pod is gone
getPod, err := kubernetesClient.CoreV1().Pods(pod.Namespace).Get(ctx, pod.Name, metav1.GetOptions{})
if err == nil {
t.Fatalf("Pod %v still exists", getPod.ObjectMeta)
}
}
+6 -6
View File
@@ -37,7 +37,7 @@ var (
)
// CleanupKubeObject deletes given kubernetes object
func CleanupKubeObject(ctx context.Context, logger *zap.Logger, kubeClient *kubernetes.Clientset, kubeobj *apiv1.ObjectReference) {
func CleanupKubeObject(ctx context.Context, logger *zap.Logger, kubeClient kubernetes.Interface, kubeobj *apiv1.ObjectReference) {
switch strings.ToLower(kubeobj.Kind) {
case "pod":
err := kubeClient.CoreV1().Pods(kubeobj.Namespace).Delete(ctx, kubeobj.Name, meta_v1.DeleteOptions{})
@@ -70,7 +70,7 @@ func CleanupKubeObject(ctx context.Context, logger *zap.Logger, kubeClient *kube
}
// CleanupDeployments deletes deployment(s) for a given instanceID
func CleanupDeployments(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error {
func CleanupDeployments(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, instanceID string, listOps meta_v1.ListOptions) error {
deploymentList, err := client.AppsV1().Deployments(meta_v1.NamespaceAll).List(ctx, listOps)
if err != nil {
return err
@@ -97,7 +97,7 @@ func CleanupDeployments(ctx context.Context, logger *zap.Logger, client *kuberne
}
// CleanupPods deletes pod(s) for a given instanceID
func CleanupPods(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error {
func CleanupPods(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, instanceID string, listOps meta_v1.ListOptions) error {
podList, err := client.CoreV1().Pods(meta_v1.NamespaceAll).List(ctx, listOps)
if err != nil {
return err
@@ -124,7 +124,7 @@ func CleanupPods(ctx context.Context, logger *zap.Logger, client *kubernetes.Cli
}
// CleanupServices deletes service(s) for a given instanceID
func CleanupServices(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error {
func CleanupServices(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, instanceID string, listOps meta_v1.ListOptions) error {
svcList, err := client.CoreV1().Services(meta_v1.NamespaceAll).List(ctx, listOps)
if err != nil {
return err
@@ -151,7 +151,7 @@ func CleanupServices(ctx context.Context, logger *zap.Logger, client *kubernetes
}
// CleanupHpa deletes horizontal pod autoscaler(s) for a given instanceID
func CleanupHpa(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, instanceID string, listOps meta_v1.ListOptions) error {
func CleanupHpa(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, instanceID string, listOps meta_v1.ListOptions) error {
hpaList, err := client.AutoscalingV1().HorizontalPodAutoscalers(meta_v1.NamespaceAll).List(ctx, listOps)
if err != nil {
return err
@@ -180,7 +180,7 @@ func CleanupHpa(ctx context.Context, logger *zap.Logger, client *kubernetes.Clie
// CleanupRoleBindings periodically lists rolebindings across all namespaces and removes Service Accounts from them or
// deletes the rolebindings completely if there are no Service Accounts in a rolebinding object.
func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client *kubernetes.Clientset, fissionClient *crd.FissionClient, functionNs, envBuilderNs string, cleanupRoleBindingInterval time.Duration) {
func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, fissionClient *crd.FissionClient, functionNs, envBuilderNs string, cleanupRoleBindingInterval time.Duration) {
for {
// some sleep before the next reaper iteration
time.Sleep(cleanupRoleBindingInterval)
+1 -1
View File
@@ -93,7 +93,7 @@ func MakeFetcherConfig(sharedMountPath string) (*Config, error) {
}, nil
}
func (cfg *Config) SetupServiceAccount(kubernetesClient *kubernetes.Clientset, namespace string, context interface{}) error {
func (cfg *Config) SetupServiceAccount(kubernetesClient kubernetes.Interface, namespace string, context interface{}) error {
_, err := utils.SetupSA(kubernetesClient, fv1.FissionFetcherSA, namespace)
if err != nil {
log.Printf("Error : %v creating %s in ns : %s for: %#v", err, fv1.FissionFetcherSA, namespace, context)
+2 -2
View File
@@ -13,7 +13,7 @@ import (
v1 "github.com/fission/fission/pkg/apis/core/v1"
)
func GetInformerFactoryByReadyPod(client *kubernetes.Clientset, namespace string, labelSelector *metav1.LabelSelector) (k8sInformers.SharedInformerFactory, error) {
func GetInformerFactoryByReadyPod(client kubernetes.Interface, namespace string, labelSelector *metav1.LabelSelector) (k8sInformers.SharedInformerFactory, error) {
informerFactory := k8sInformers.NewSharedInformerFactoryWithOptions(client, 0,
k8sInformers.WithNamespace(namespace),
k8sInformers.WithTweakListOptions(func(options *metav1.ListOptions) {
@@ -23,7 +23,7 @@ func GetInformerFactoryByReadyPod(client *kubernetes.Clientset, namespace string
return informerFactory, nil
}
func GetInformerFactoryByExecutor(client *kubernetes.Clientset, executorType v1.ExecutorType, defaultResync time.Duration) (k8sInformers.SharedInformerFactory, error) {
func GetInformerFactoryByExecutor(client kubernetes.Interface, 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
+5 -5
View File
@@ -50,7 +50,7 @@ func MakeSAObj(sa, ns string) *apiv1.ServiceAccount {
}
// SetupSA checks if a service account is present in the namespace, if not creates it.
func SetupSA(k8sClient *kubernetes.Clientset, sa, ns string) (*apiv1.ServiceAccount, error) {
func SetupSA(k8sClient kubernetes.Interface, sa, ns string) (*apiv1.ServiceAccount, error) {
saObj, err := k8sClient.CoreV1().ServiceAccounts(ns).Get(context.TODO(), sa, metav1.GetOptions{})
if err == nil {
return saObj, nil
@@ -105,7 +105,7 @@ type PatchSpec struct {
}
// AddSaToRoleBindingWithRetries adds a service account to a rolebinding object. IT retries on already exists and conflict errors.
func AddSaToRoleBindingWithRetries(ctx context.Context, logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind string) (err error) {
func AddSaToRoleBindingWithRetries(ctx context.Context, logger *zap.Logger, k8sClient kubernetes.Interface, roleBinding, roleBindingNs, sa, saNamespace, role, roleKind string) (err error) {
patch := PatchSpec{}
patch.Op = "add"
patch.Path = "/subjects/-"
@@ -177,7 +177,7 @@ func AddSaToRoleBindingWithRetries(ctx context.Context, logger *zap.Logger, k8sC
// RemoveSAFromRoleBindingWithRetries removes an SA from the rolebinding passed as parameter. If this is the only SA in
// the rolebinding, then it deletes the rolebinding object.
func RemoveSAFromRoleBindingWithRetries(ctx context.Context, logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs string, saToRemove map[string]bool) (err error) {
func RemoveSAFromRoleBindingWithRetries(ctx context.Context, logger *zap.Logger, k8sClient kubernetes.Interface, roleBinding, roleBindingNs string, saToRemove map[string]bool) (err error) {
for i := 0; i < maxRetries; i++ {
rbObj, err := k8sClient.RbacV1().RoleBindings(roleBindingNs).Get(
ctx,
@@ -237,7 +237,7 @@ func RemoveSAFromRoleBindingWithRetries(ctx context.Context, logger *zap.Logger,
// SetupRoleBinding adds a role to a service account if the rolebinding object is already present in the namespace.
// if not, it creates a rolebinding object granting the role to the SA in the namespace.
func SetupRoleBinding(ctx context.Context, logger *zap.Logger, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs, role, roleKind, sa, saNamespace string) error {
func SetupRoleBinding(ctx context.Context, logger *zap.Logger, k8sClient kubernetes.Interface, roleBinding, roleBindingNs, role, roleKind, sa, saNamespace string) error {
// get the role binding object
rbObj, err := k8sClient.RbacV1().RoleBindings(roleBindingNs).Get(
ctx,
@@ -283,7 +283,7 @@ func SetupRoleBinding(ctx context.Context, logger *zap.Logger, k8sClient *kubern
// DeleteRoleBinding deletes a rolebinding object. if k8s throws an error that the rolebinding is not there, it just
// returns silently.
func DeleteRoleBinding(ctx context.Context, k8sClient *kubernetes.Clientset, roleBinding, roleBindingNs string) error {
func DeleteRoleBinding(ctx context.Context, k8sClient kubernetes.Interface, roleBinding, roleBindingNs string) error {
// if deleteRoleBinding is invoked by 2 fission services at the same time for the same rolebinding,
// the first call will succeed while the 2nd will fail with isNotFound. but we dont want to error out then.
err := k8sClient.RbacV1().RoleBindings(roleBindingNs).Delete(ctx, roleBinding, metav1.DeleteOptions{})