Allow empty namespace for fission function and builder (#2621)
Currently, we create Fission resources in the default namespace, function-related resources are created in the fission-function namespace, whereas builder resources are created in the fission-builder namespace. This causes confusion for a lot of users. In this fix, we allow the user to set the function and builder namespace empty so that function and builder resources are created in the same namespace as the function resource always. If the user desires older behaviour they can functionNamespace and builderNamespace the same previous before the upgrade. * use default namespace for fission function and builder * support for existing fission namespaces * Replace builder and function namespace with template * Fix namespace creation template Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
@@ -255,7 +255,7 @@ func (executor *Executor) getFunctionServiceFromCache(ctx context.Context, fn *f
|
||||
|
||||
// StartExecutor Starts executor and the executor components such as Poolmgr,
|
||||
// deploymgr and potential future executor types
|
||||
func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace string, envBuilderNamespace string, port int) error {
|
||||
func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error {
|
||||
fissionClient, kubernetesClient, _, metricsClient, err := crd.MakeFissionClient()
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "failed to get kubernetes client")
|
||||
@@ -306,7 +306,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st
|
||||
gpm, err := poolmgr.MakeGenericPoolManager(ctx,
|
||||
logger,
|
||||
fissionClient, kubernetesClient, metricsClient,
|
||||
functionNamespace, fetcherConfig, executorInstanceID,
|
||||
fetcherConfig, executorInstanceID,
|
||||
funcInformer, pkgInformer, envInformer,
|
||||
gpmPodInformer, gpmRsInformer, podSpecPatch)
|
||||
if err != nil {
|
||||
@@ -322,7 +322,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st
|
||||
ndm, err := newdeploy.MakeNewDeploy(ctx,
|
||||
logger,
|
||||
fissionClient, kubernetesClient,
|
||||
functionNamespace, fetcherConfig, executorInstanceID,
|
||||
fetcherConfig, executorInstanceID,
|
||||
funcInformer, envInformer,
|
||||
ndmDeplInformer, ndmSvcInformer, podSpecPatch)
|
||||
if err != nil {
|
||||
@@ -338,7 +338,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st
|
||||
cnm, err := container.MakeContainer(
|
||||
ctx, logger,
|
||||
fissionClient, kubernetesClient,
|
||||
functionNamespace, executorInstanceID, funcInformer,
|
||||
executorInstanceID, funcInformer,
|
||||
cnmDeplInformer, cnmSvcInformer)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "container manager creation failed")
|
||||
@@ -408,7 +408,8 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
go reaper.CleanupRoleBindings(ctx, logger, kubernetesClient, fissionClient, functionNamespace, envBuilderNamespace, time.Minute*30)
|
||||
|
||||
go reaper.CleanupRoleBindings(ctx, logger, kubernetesClient, fissionClient, time.Minute*30)
|
||||
go metrics.ServeMetrics(ctx, logger)
|
||||
go api.Serve(ctx, port)
|
||||
|
||||
|
||||
@@ -174,7 +174,7 @@ func TestExecutor(t *testing.T) {
|
||||
|
||||
// create poolmgr
|
||||
port := 9999
|
||||
err = StartExecutor(ctx, logger, functionNs, "fission-builder", port)
|
||||
err = StartExecutor(ctx, logger, port)
|
||||
if err != nil {
|
||||
log.Panicf("failed to start poolmgr: %v", err)
|
||||
}
|
||||
|
||||
@@ -57,7 +57,9 @@ import (
|
||||
otelUtils "github.com/fission/fission/pkg/utils/otel"
|
||||
)
|
||||
|
||||
var _ executortype.ExecutorType = &Container{}
|
||||
var (
|
||||
_ executortype.ExecutorType = &Container{}
|
||||
)
|
||||
|
||||
type (
|
||||
// Container represents an executor type
|
||||
@@ -67,10 +69,10 @@ type (
|
||||
kubernetesClient kubernetes.Interface
|
||||
fissionClient versioned.Interface
|
||||
instanceID string
|
||||
nsResolver *utils.NamespaceResolver
|
||||
// fetcherConfig *fetcherConfig.Config
|
||||
|
||||
runtimeImagePullPolicy apiv1.PullPolicy
|
||||
namespace string
|
||||
useIstio bool
|
||||
|
||||
fsCache *fscache.FunctionServiceCache // cache funcSvc's by function, address and pod name
|
||||
@@ -96,7 +98,6 @@ func MakeContainer(
|
||||
logger *zap.Logger,
|
||||
fissionClient versioned.Interface,
|
||||
kubernetesClient kubernetes.Interface,
|
||||
namespace string,
|
||||
instanceID string,
|
||||
funcInformer map[string]finformerv1.FunctionInformer,
|
||||
deplInformer appsinformers.DeploymentInformer,
|
||||
@@ -117,8 +118,8 @@ func MakeContainer(
|
||||
fissionClient: fissionClient,
|
||||
kubernetesClient: kubernetesClient,
|
||||
instanceID: instanceID,
|
||||
nsResolver: utils.DefaultNSResolver(),
|
||||
|
||||
namespace: namespace,
|
||||
fsCache: fscache.MakeFunctionServiceCache(logger),
|
||||
throttler: throttler.MakeThrottler(1 * time.Minute),
|
||||
|
||||
@@ -129,6 +130,7 @@ func MakeContainer(
|
||||
objectReaperIntervalSecond: time.Duration(executorUtils.GetObjectReaperInterval(logger, fv1.ExecutorTypeContainer, 5)) * time.Second,
|
||||
hpaops: hpautils.NewHpaOperations(logger, kubernetesClient, instanceID),
|
||||
}
|
||||
|
||||
caaf.deplLister = deplInformer.Lister()
|
||||
caaf.deplListerSynced = deplInformer.Informer().HasSynced
|
||||
|
||||
@@ -382,10 +384,7 @@ func (caaf *Container) fnCreate(ctx context.Context, fn *fv1.Function) (*fscache
|
||||
|
||||
// to support backward compatibility, if the function was created in default ns, we fall back to creating the
|
||||
// deployment of the function in fission-function ns
|
||||
ns := caaf.namespace
|
||||
if fn.ObjectMeta.Namespace != metav1.NamespaceDefault {
|
||||
ns = fn.ObjectMeta.Namespace
|
||||
}
|
||||
ns := caaf.nsResolver.GetFunctionNS(fn.ObjectMeta.Namespace)
|
||||
|
||||
// Envoy(istio-proxy) returns 404 directly before istio pilot
|
||||
// propagates latest Envoy-specific configuration.
|
||||
@@ -499,10 +498,7 @@ func (caaf *Container) updateFunction(ctx context.Context, oldFn *fv1.Function,
|
||||
if !reflect.DeepEqual(oldFn.Spec.InvokeStrategy, newFn.Spec.InvokeStrategy) {
|
||||
// to support backward compatibility, if the function was created in default ns, we fall back to creating the
|
||||
// deployment of the function in fission-function ns, so cleaning up resources there
|
||||
ns := caaf.namespace
|
||||
if newFn.ObjectMeta.Namespace != metav1.NamespaceDefault {
|
||||
ns = newFn.ObjectMeta.Namespace
|
||||
}
|
||||
ns := caaf.nsResolver.GetFunctionNS(newFn.ObjectMeta.Namespace)
|
||||
|
||||
fsvc, err := caaf.fsCache.GetByFunctionUID(newFn.ObjectMeta.UID)
|
||||
if err != nil {
|
||||
@@ -598,10 +594,7 @@ func (caaf *Container) updateFuncDeployment(ctx context.Context, fn *fv1.Functio
|
||||
|
||||
// to support backward compatibility, if the function was created in default ns, we fall back to creating the
|
||||
// deployment of the function in fission-function ns
|
||||
ns := caaf.namespace
|
||||
if fn.ObjectMeta.Namespace != metav1.NamespaceDefault {
|
||||
ns = fn.ObjectMeta.Namespace
|
||||
}
|
||||
ns := caaf.nsResolver.GetFunctionNS(fn.ObjectMeta.Namespace)
|
||||
|
||||
existingDepl, err := caaf.kubernetesClient.AppsV1().Deployments(ns).Get(ctx, fnObjName, metav1.GetOptions{})
|
||||
if err != nil {
|
||||
@@ -651,10 +644,7 @@ func (caaf *Container) fnDelete(ctx context.Context, fn *fv1.Function) error {
|
||||
|
||||
// to support backward compatibility, if the function was created in default ns, we fall back to creating the
|
||||
// deployment of the function in fission-function ns, so cleaning up resources there
|
||||
ns := caaf.namespace
|
||||
if fn.ObjectMeta.Namespace != metav1.NamespaceDefault {
|
||||
ns = fn.ObjectMeta.Namespace
|
||||
}
|
||||
ns := caaf.nsResolver.GetFunctionNS(fn.ObjectMeta.Namespace)
|
||||
|
||||
err = caaf.cleanupContainer(ctx, ns, objName)
|
||||
multierr = multierror.Append(multierr, err)
|
||||
|
||||
@@ -59,7 +59,9 @@ import (
|
||||
otelUtils "github.com/fission/fission/pkg/utils/otel"
|
||||
)
|
||||
|
||||
var _ executortype.ExecutorType = &NewDeploy{}
|
||||
var (
|
||||
_ executortype.ExecutorType = &NewDeploy{}
|
||||
)
|
||||
|
||||
type (
|
||||
// NewDeploy represents an ExecutorType
|
||||
@@ -70,9 +72,9 @@ type (
|
||||
fissionClient versioned.Interface
|
||||
instanceID string
|
||||
fetcherConfig *fetcherConfig.Config
|
||||
nsResolver *utils.NamespaceResolver
|
||||
|
||||
runtimeImagePullPolicy apiv1.PullPolicy
|
||||
namespace string
|
||||
useIstio bool
|
||||
|
||||
fsCache *fscache.FunctionServiceCache // cache funcSvc's by function, address and pod name
|
||||
@@ -100,7 +102,6 @@ func MakeNewDeploy(
|
||||
logger *zap.Logger,
|
||||
fissionClient versioned.Interface,
|
||||
kubernetesClient kubernetes.Interface,
|
||||
namespace string,
|
||||
fetcherConfig *fetcherConfig.Config,
|
||||
instanceID string,
|
||||
funcInformer map[string]finformerv1.FunctionInformer,
|
||||
@@ -124,10 +125,9 @@ func MakeNewDeploy(
|
||||
fissionClient: fissionClient,
|
||||
kubernetesClient: kubernetesClient,
|
||||
instanceID: instanceID,
|
||||
|
||||
namespace: namespace,
|
||||
fsCache: fscache.MakeFunctionServiceCache(logger),
|
||||
throttler: throttler.MakeThrottler(1 * time.Minute),
|
||||
fsCache: fscache.MakeFunctionServiceCache(logger),
|
||||
throttler: throttler.MakeThrottler(1 * time.Minute),
|
||||
nsResolver: utils.DefaultNSResolver(),
|
||||
|
||||
fetcherConfig: fetcherConfig,
|
||||
runtimeImagePullPolicy: utils.GetImagePullPolicy(os.Getenv("RUNTIME_IMAGE_PULL_POLICY")),
|
||||
@@ -428,10 +428,7 @@ func (deploy *NewDeploy) fnCreate(ctx context.Context, fn *fv1.Function) (*fscac
|
||||
|
||||
// to support backward compatibility, if the function was created in default ns, we fall back to creating the
|
||||
// deployment of the function in fission-function ns
|
||||
ns := deploy.namespace
|
||||
if fn.ObjectMeta.Namespace != metav1.NamespaceDefault {
|
||||
ns = fn.ObjectMeta.Namespace
|
||||
}
|
||||
ns := deploy.nsResolver.GetFunctionNS(fn.ObjectMeta.Namespace)
|
||||
|
||||
// Envoy(istio-proxy) returns 404 directly before istio pilot
|
||||
// propagates latest Envoy-specific configuration.
|
||||
@@ -549,10 +546,7 @@ func (deploy *NewDeploy) updateFunction(ctx context.Context, oldFn *fv1.Function
|
||||
|
||||
// to support backward compatibility, if the function was created in default ns, we fall back to creating the
|
||||
// deployment of the function in fission-function ns, so cleaning up resources there
|
||||
ns := deploy.namespace
|
||||
if newFn.ObjectMeta.Namespace != metav1.NamespaceDefault {
|
||||
ns = newFn.ObjectMeta.Namespace
|
||||
}
|
||||
ns := deploy.nsResolver.GetFunctionNS(newFn.ObjectMeta.Namespace)
|
||||
|
||||
fsvc, err := deploy.fsCache.GetByFunctionUID(newFn.ObjectMeta.UID)
|
||||
if err != nil {
|
||||
@@ -654,10 +648,7 @@ func (deploy *NewDeploy) updateFuncDeployment(ctx context.Context, fn *fv1.Funct
|
||||
|
||||
// to support backward compatibility, if the function was created in default ns, we fall back to creating the
|
||||
// deployment of the function in fission-function ns
|
||||
ns := deploy.namespace
|
||||
if fn.ObjectMeta.Namespace != metav1.NamespaceDefault {
|
||||
ns = fn.ObjectMeta.Namespace
|
||||
}
|
||||
ns := deploy.nsResolver.GetFunctionNS(fn.ObjectMeta.Namespace)
|
||||
|
||||
existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(ns).Get(ctx, fnObjName, metav1.GetOptions{})
|
||||
if err != nil {
|
||||
@@ -708,10 +699,7 @@ func (deploy *NewDeploy) fnDelete(ctx context.Context, fn *fv1.Function) error {
|
||||
|
||||
// to support backward compatibility, if the function was created in default ns, we fall back to creating the
|
||||
// deployment of the function in fission-function ns, so cleaning up resources there
|
||||
ns := deploy.namespace
|
||||
if fn.ObjectMeta.Namespace != metav1.NamespaceDefault {
|
||||
ns = fn.ObjectMeta.Namespace
|
||||
}
|
||||
ns := deploy.nsResolver.GetFunctionNS(fn.ObjectMeta.Namespace)
|
||||
|
||||
err = deploy.cleanupNewdeploy(ctx, ns, objName)
|
||||
multierr = multierror.Append(multierr, err)
|
||||
|
||||
@@ -30,6 +30,7 @@ import (
|
||||
const (
|
||||
defaultNamespace string = "default"
|
||||
functionNamespace string = "fission-function"
|
||||
builderNamespace string = "fission-builder"
|
||||
envName string = "newdeploy-test-env"
|
||||
functionName string = "newdeploy-test-func"
|
||||
configmapName string = "newdeploy-test-configmap"
|
||||
@@ -80,7 +81,7 @@ func TestRefreshFuncPods(t *testing.T) {
|
||||
t.Fatalf("Error creating fetcher config: %s", err)
|
||||
}
|
||||
|
||||
executor, err := MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, functionNamespace, fetcherConfig, "test",
|
||||
executor, err := MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, fetcherConfig, "test",
|
||||
funcInformer, envInformer, deployInformer, svcInformer, podSpecPatch)
|
||||
if err != nil {
|
||||
t.Fatalf("new deploy manager creation failed: %s", err)
|
||||
@@ -88,6 +89,12 @@ func TestRefreshFuncPods(t *testing.T) {
|
||||
|
||||
ndm := executor.(*NewDeploy)
|
||||
|
||||
nsResolver := utils.NamespaceResolver{
|
||||
FunctionNamespace: functionNamespace,
|
||||
BuiderNamespace: builderNamespace,
|
||||
}
|
||||
ndm.nsResolver = &nsResolver
|
||||
|
||||
go ndm.Run(ctx)
|
||||
t.Log("New deploy manager started")
|
||||
|
||||
|
||||
@@ -62,7 +62,6 @@ type (
|
||||
env *fv1.Environment
|
||||
deployment *appsv1.Deployment // kubernetes deployment
|
||||
namespace string // namespace to keep our resources
|
||||
functionNamespace string // fallback namespace for fission functions
|
||||
podReadyTimeout time.Duration // timeout for generic pods to become ready
|
||||
fsCache *fscache.FunctionServiceCache // cache funcSvc's by function, address and podname
|
||||
useSvc bool // create k8s service for specialized pods
|
||||
@@ -92,7 +91,6 @@ func MakeGenericPool(
|
||||
metricsClient metricsclient.Interface,
|
||||
env *fv1.Environment,
|
||||
namespace string,
|
||||
functionNamespace string,
|
||||
fsCache *fscache.FunctionServiceCache,
|
||||
fetcherConfig *fetcherConfig.Config,
|
||||
instanceID string,
|
||||
@@ -122,7 +120,6 @@ func MakeGenericPool(
|
||||
kubernetesClient: kubernetesClient,
|
||||
metricsClient: metricsClient,
|
||||
namespace: namespace,
|
||||
functionNamespace: functionNamespace,
|
||||
podReadyTimeout: podReadyTimeout,
|
||||
fsCache: fsCache,
|
||||
fetcherConfig: fetcherConfig,
|
||||
@@ -446,7 +443,7 @@ func (gp *GenericPool) createSvc(ctx context.Context, name string, labels map[st
|
||||
}
|
||||
|
||||
func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscache.FuncSvc, error) {
|
||||
logger := otelUtils.LoggerWithTraceID(ctx, gp.logger).With(zap.String("function", fn.ObjectMeta.Name), zap.String("functionNamespace", fn.ObjectMeta.Namespace),
|
||||
logger := otelUtils.LoggerWithTraceID(ctx, gp.logger).With(zap.String("function", fn.ObjectMeta.Name), zap.String("namespace", fn.ObjectMeta.Namespace),
|
||||
zap.String("env", fn.Spec.Environment.Name), zap.String("envNamespace", fn.Spec.Environment.Namespace))
|
||||
|
||||
logger.Info("choosing pod from pool")
|
||||
|
||||
@@ -58,7 +58,9 @@ import (
|
||||
otelUtils "github.com/fission/fission/pkg/utils/otel"
|
||||
)
|
||||
|
||||
var _ executortype.ExecutorType = &GenericPoolManager{}
|
||||
var (
|
||||
_ executortype.ExecutorType = &GenericPoolManager{}
|
||||
)
|
||||
|
||||
type requestType int
|
||||
|
||||
@@ -74,7 +76,7 @@ type (
|
||||
pools map[string]*GenericPool
|
||||
kubernetesClient kubernetes.Interface
|
||||
metricsClient metricsclient.Interface
|
||||
namespace string
|
||||
nsResolver *utils.NamespaceResolver
|
||||
|
||||
fissionClient versioned.Interface
|
||||
functionEnv *cache.Cache
|
||||
@@ -116,7 +118,6 @@ func MakeGenericPoolManager(ctx context.Context,
|
||||
fissionClient versioned.Interface,
|
||||
kubernetesClient kubernetes.Interface,
|
||||
metricsClient metricsclient.Interface,
|
||||
functionNamespace string,
|
||||
fetcherConfig *fetcherConfig.Config,
|
||||
instanceID string,
|
||||
funcInformer map[string]finformerv1.FunctionInformer,
|
||||
@@ -138,15 +139,15 @@ func MakeGenericPoolManager(ctx context.Context,
|
||||
enableIstio = istio
|
||||
}
|
||||
|
||||
poolPodC := NewPoolPodController(ctx, gpmLogger, kubernetesClient, functionNamespace,
|
||||
poolPodC := NewPoolPodController(ctx, gpmLogger, kubernetesClient,
|
||||
enableIstio, funcInformer, pkgInformer, envInformer, rsInformer, podInformer)
|
||||
|
||||
gpm := &GenericPoolManager{
|
||||
logger: gpmLogger,
|
||||
pools: make(map[string]*GenericPool),
|
||||
kubernetesClient: kubernetesClient,
|
||||
nsResolver: utils.DefaultNSResolver(),
|
||||
metricsClient: metricsClient,
|
||||
namespace: functionNamespace,
|
||||
fissionClient: fissionClient,
|
||||
functionEnv: cache.MakeCache(10*time.Second, 0),
|
||||
fsCache: fscache.MakeFunctionServiceCache(gpmLogger),
|
||||
@@ -162,6 +163,8 @@ func MakeGenericPoolManager(ctx context.Context,
|
||||
gpm.podLister = podInformer.Lister()
|
||||
gpm.podListerSynced = podInformer.Informer().HasSynced
|
||||
|
||||
gpm.logger.Debug("inside MakeGenericPoolManager")
|
||||
|
||||
return gpm, nil
|
||||
}
|
||||
|
||||
@@ -197,7 +200,7 @@ func (gpm *GenericPoolManager) GetFuncSvc(ctx context.Context, fn *fv1.Function)
|
||||
}
|
||||
|
||||
if created {
|
||||
logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.namespace))
|
||||
logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.nsResolver.ResolveNamespace(gpm.nsResolver.FunctionNamespace)))
|
||||
}
|
||||
|
||||
// from GenericPool -> get one function container
|
||||
@@ -277,7 +280,7 @@ func (gpm *GenericPoolManager) RefreshFuncPods(ctx context.Context, logger *zap.
|
||||
}
|
||||
|
||||
if created {
|
||||
gpm.logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.namespace))
|
||||
gpm.logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.nsResolver.ResolveNamespace(gpm.nsResolver.FunctionNamespace)))
|
||||
}
|
||||
|
||||
funcSvc, err := gp.fsCache.GetByFunction(&f.ObjectMeta)
|
||||
@@ -333,7 +336,7 @@ func (gpm *GenericPoolManager) AdoptExistingResources(ctx context.Context) {
|
||||
gpm.logger.Error("adopt 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))
|
||||
gpm.logger.Info("created pool for the environment", zap.String("env", env.ObjectMeta.Name), zap.String("namespace", gpm.nsResolver.ResolveNamespace(gpm.nsResolver.FunctionNamespace)))
|
||||
}
|
||||
}()
|
||||
}
|
||||
@@ -482,12 +485,9 @@ func (gpm *GenericPoolManager) service() {
|
||||
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.
|
||||
ns := gpm.namespace
|
||||
if req.env.ObjectMeta.Namespace != metav1.NamespaceDefault {
|
||||
ns = req.env.ObjectMeta.Namespace
|
||||
}
|
||||
ns := gpm.nsResolver.GetFunctionNS(req.env.ObjectMeta.Namespace)
|
||||
pool = MakeGenericPool(gpm.logger, gpm.fissionClient, gpm.kubernetesClient,
|
||||
gpm.metricsClient, req.env, ns, gpm.namespace, gpm.fsCache,
|
||||
gpm.metricsClient, req.env, ns, gpm.fsCache,
|
||||
gpm.fetcherConfig, gpm.instanceID, gpm.enableIstio, gpm.podSpecPatch)
|
||||
err = pool.setup(req.ctx)
|
||||
if err != nil {
|
||||
|
||||
@@ -40,14 +40,15 @@ import (
|
||||
"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"
|
||||
"github.com/fission/fission/pkg/utils"
|
||||
)
|
||||
|
||||
type (
|
||||
PoolPodController struct {
|
||||
logger *zap.Logger
|
||||
kubernetesClient kubernetes.Interface
|
||||
namespace string
|
||||
enableIstio bool
|
||||
nsResolver *utils.NamespaceResolver
|
||||
|
||||
envLister map[string]flisterv1.EnvironmentLister
|
||||
envListerSynced map[string]k8sCache.InformerSynced
|
||||
@@ -69,7 +70,6 @@ type (
|
||||
|
||||
func NewPoolPodController(ctx context.Context, logger *zap.Logger,
|
||||
kubernetesClient kubernetes.Interface,
|
||||
namespace string,
|
||||
enableIstio bool,
|
||||
funcInformer map[string]finformerv1.FunctionInformer,
|
||||
pkgInformer map[string]finformerv1.PackageInformer,
|
||||
@@ -79,8 +79,8 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger,
|
||||
logger = logger.Named("pool_pod_controller")
|
||||
p := &PoolPodController{
|
||||
logger: logger,
|
||||
nsResolver: utils.DefaultNSResolver(),
|
||||
kubernetesClient: kubernetesClient,
|
||||
namespace: namespace,
|
||||
enableIstio: enableIstio,
|
||||
envLister: make(map[string]flisterv1.EnvironmentLister, 0),
|
||||
envListerSynced: make(map[string]k8sCache.InformerSynced, 0),
|
||||
@@ -89,10 +89,10 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger,
|
||||
spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"),
|
||||
}
|
||||
for _, informer := range funcInformer {
|
||||
informer.Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace, p.enableIstio))
|
||||
informer.Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace), p.enableIstio))
|
||||
}
|
||||
for _, informer := range pkgInformer {
|
||||
informer.Informer().AddEventHandler(PackageEventHandlers(ctx, p.logger, p.kubernetesClient, p.namespace))
|
||||
informer.Informer().AddEventHandler(PackageEventHandlers(ctx, p.logger, p.kubernetesClient, p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace)))
|
||||
}
|
||||
for ns, informer := range envInformer {
|
||||
informer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
|
||||
@@ -373,7 +373,7 @@ 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.gpm.namespace).List(labels.SelectorFromSet(specializePodLables))
|
||||
specializedPods, err := p.podLister.Pods(p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace)).List(labels.SelectorFromSet(specializePodLables))
|
||||
if err != nil {
|
||||
p.logger.Error("failed to list specialized pods", zap.Error(err))
|
||||
p.envDeleteQueue.Forget(obj)
|
||||
|
||||
@@ -68,8 +68,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
|
||||
gpmPodInformer := gpmInformerFactory.Core().V1().Pods()
|
||||
gpmRsInformer := gpmInformerFactory.Apps().V1().ReplicaSets()
|
||||
|
||||
fnNamespace := "fission-function"
|
||||
ppc := NewPoolPodController(ctx, logger, kubernetesClient, fnNamespace, false,
|
||||
ppc := NewPoolPodController(ctx, logger, kubernetesClient, false,
|
||||
funcInformer,
|
||||
pkgInformer,
|
||||
envInformer,
|
||||
@@ -85,7 +84,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
|
||||
executor, err := MakeGenericPoolManager(ctx,
|
||||
logger,
|
||||
fissionClient, kubernetesClient, metricsClient,
|
||||
fnNamespace, fetcherConfig, executorInstanceID,
|
||||
fetcherConfig, executorInstanceID,
|
||||
funcInformer, pkgInformer, envInformer,
|
||||
gpmPodInformer, gpmRsInformer, nil)
|
||||
if err != nil {
|
||||
|
||||
@@ -184,7 +184,8 @@ func CleanupHpa(ctx context.Context, logger *zap.Logger, client kubernetes.Inter
|
||||
|
||||
// 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.Interface, fissionClient versioned.Interface, functionNs, envBuilderNs string, cleanupRoleBindingInterval time.Duration) {
|
||||
func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client kubernetes.Interface, fissionClient versioned.Interface, cleanupRoleBindingInterval time.Duration) {
|
||||
nsResolver := utils.DefaultNSResolver()
|
||||
for {
|
||||
// some sleep before the next reaper iteration
|
||||
time.Sleep(cleanupRoleBindingInterval)
|
||||
@@ -238,8 +239,8 @@ func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client kuberne
|
||||
// so now we need to look for the objects in default namespace.
|
||||
saNs := subj.Namespace
|
||||
isInReservedNS := false
|
||||
if subj.Namespace == functionNs ||
|
||||
subj.Namespace == envBuilderNs {
|
||||
if subj.Namespace == nsResolver.FunctionNamespace ||
|
||||
subj.Namespace == nsResolver.BuiderNamespace {
|
||||
saNs = metav1.NamespaceDefault
|
||||
isInReservedNS = true
|
||||
}
|
||||
@@ -249,7 +250,8 @@ func CleanupRoleBindings(ctx context.Context, logger *zap.Logger, client kuberne
|
||||
for _, fn := range funcList.Items {
|
||||
if fn.Spec.Environment.Namespace == saNs ||
|
||||
// For the case that the environment is created in the reserved namespace.
|
||||
(isInReservedNS && (fn.Spec.Environment.Namespace == functionNs || fn.Spec.Environment.Namespace == envBuilderNs)) {
|
||||
(isInReservedNS && (fn.Spec.Environment.Namespace == nsResolver.FunctionNamespace ||
|
||||
fn.Spec.Environment.Namespace == nsResolver.BuiderNamespace)) {
|
||||
funcEnvReference = true
|
||||
break
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user