Executor user informer factory in executors in place of informers (#2666)

* Use informerfactory across executor
* Run function informer for poolpodcontroller if istio enabled
* Use same namespace for secret as keda mqtriggers

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
Sanket Sudake
2022-12-11 20:41:12 +05:30
committed by GitHub
parent d16de59e9f
commit 3ae1742953
10 changed files with 60 additions and 115 deletions
+6 -1
View File
@@ -33,6 +33,7 @@ import (
metricsclient "k8s.io/metrics/pkg/client/clientset/versioned"
"github.com/fission/fission/pkg/generated/clientset/versioned"
"github.com/fission/fission/pkg/utils"
)
// GetKubernetesClient gets a kubernetes client using the kubeconfig file at the
@@ -91,9 +92,13 @@ func MakeFissionClient() (versioned.Interface, kubernetes.Interface, apiextensio
// WaitForCRDs does a timeout to check if CRDs have been installed
func WaitForCRDs(ctx context.Context, logger *zap.Logger, fissionClient versioned.Interface) error {
logger.Info("Waiting for CRDs to be installed")
defaultNs := utils.DefaultNSResolver().DefaultNamespace
if defaultNs == "" {
defaultNs = metav1.NamespaceDefault
}
start := time.Now()
for {
fi := fissionClient.CoreV1().Functions(metav1.NamespaceDefault)
fi := fissionClient.CoreV1().Functions(defaultNs)
_, err := fi.List(ctx, metav1.ListOptions{})
if err != nil {
time.Sleep(100 * time.Millisecond)
+11 -22
View File
@@ -42,7 +42,6 @@ import (
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
"github.com/fission/fission/pkg/generated/clientset/versioned"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/metrics"
otelUtils "github.com/fission/fission/pkg/utils/otel"
@@ -276,13 +275,9 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error {
logger.Info("Starting executor", zap.String("instanceID", executorInstanceID))
funcInformer := make(map[string]finformerv1.FunctionInformer, 0)
envInformer := make(map[string]finformerv1.EnvironmentInformer, 0)
finformerFactory := make(map[string]genInformer.SharedInformerFactory, 0)
for _, ns := range utils.DefaultNSResolver().FissionResourceNS {
factory := genInformer.NewFilteredSharedInformerFactory(fissionClient, time.Minute*30, ns, nil)
funcInformer[ns] = factory.Core().V1().Functions()
envInformer[ns] = factory.Core().V1().Environments()
finformerFactory[ns] = genInformer.NewFilteredSharedInformerFactory(fissionClient, time.Minute*30, ns, nil)
}
executorLabel, err := utils.GetInformerLabelByExecutor(fv1.ExecutorTypePoolmgr)
@@ -294,7 +289,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error {
logger,
fissionClient, kubernetesClient, metricsClient,
fetcherConfig, executorInstanceID,
funcInformer, envInformer,
finformerFactory,
gpmInformerFactory, podSpecPatch)
if err != nil {
return errors.Wrap(err, "pool manager creation failed")
@@ -309,7 +304,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error {
logger,
fissionClient, kubernetesClient,
fetcherConfig, executorInstanceID,
funcInformer, envInformer,
finformerFactory,
ndmInformerFactory, podSpecPatch)
if err != nil {
return errors.Wrap(err, "new deploy manager creation failed")
@@ -323,7 +318,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error {
cnm, err := container.MakeContainer(
ctx, logger,
fissionClient, kubernetesClient,
executorInstanceID, funcInformer,
executorInstanceID, finformerFactory,
cnmInformerFactory)
if err != nil {
return errors.Wrap(err, "container manager creation failed")
@@ -356,29 +351,23 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, port int) error {
cms := cms.MakeConfigSecretController(ctx, logger, fissionClient, kubernetesClient, executorTypes, configMapInformer, secretInformer)
fissionInformers := make([]k8sCache.SharedIndexInformer, 0)
for _, informer := range funcInformer {
fissionInformers = append(fissionInformers, informer.Informer())
}
for _, informer := range envInformer {
fissionInformers = append(fissionInformers, informer.Informer())
}
for _, informer := range configMapInformer {
fissionInformers = append(fissionInformers, informer)
}
for _, informer := range secretInformer {
fissionInformers = append(fissionInformers, informer)
}
for _, factory := range finformerFactory {
factory.Start(ctx.Done())
}
for _, informerFactory := range gpmInformerFactory {
fissionInformers = append(fissionInformers, informerFactory.Core().V1().Pods().Informer())
fissionInformers = append(fissionInformers, informerFactory.Apps().V1().ReplicaSets().Informer())
informerFactory.Start(ctx.Done())
}
for _, informerFactory := range ndmInformerFactory {
fissionInformers = append(fissionInformers, informerFactory.Apps().V1().Deployments().Informer())
fissionInformers = append(fissionInformers, informerFactory.Core().V1().Services().Informer())
informerFactory.Start(ctx.Done())
}
for _, informerFactory := range cnmInformerFactory {
fissionInformers = append(fissionInformers, informerFactory.Apps().V1().Deployments().Informer())
fissionInformers = append(fissionInformers, informerFactory.Core().V1().Services().Informer())
informerFactory.Start(ctx.Done())
}
api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes,
@@ -49,7 +49,7 @@ import (
executorUtils "github.com/fission/fission/pkg/executor/util"
hpautils "github.com/fission/fission/pkg/executor/util/hpa"
"github.com/fission/fission/pkg/generated/clientset/versioned"
finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
"github.com/fission/fission/pkg/throttler"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/maps"
@@ -98,7 +98,7 @@ func MakeContainer(
fissionClient versioned.Interface,
kubernetesClient kubernetes.Interface,
instanceID string,
funcInformer map[string]finformerv1.FunctionInformer,
finformerFactory map[string]genInformer.SharedInformerFactory,
cnmInformerFactory map[string]k8sInformers.SharedInformerFactory,
) (executortype.ExecutorType, error) {
enableIstio := false
@@ -139,8 +139,8 @@ func MakeContainer(
caaf.svcLister[ns] = informerFactory.Core().V1().Services().Lister()
caaf.svcListerSynced[ns] = informerFactory.Core().V1().Services().Informer().HasSynced
}
for _, informer := range funcInformer {
informer.Informer().AddEventHandler(caaf.FuncInformerHandler(ctx))
for _, factory := range finformerFactory {
factory.Core().V1().Functions().Informer().AddEventHandler(caaf.FuncInformerHandler(ctx))
}
return caaf, nil
}
@@ -51,7 +51,7 @@ import (
hpautils "github.com/fission/fission/pkg/executor/util/hpa"
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"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
"github.com/fission/fission/pkg/throttler"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/maps"
@@ -103,8 +103,7 @@ func MakeNewDeploy(
kubernetesClient kubernetes.Interface,
fetcherConfig *fetcherConfig.Config,
instanceID string,
funcInformer map[string]finformerv1.FunctionInformer,
envInformer map[string]finformerv1.EnvironmentInformer,
finformerFactory map[string]genInformer.SharedInformerFactory,
ndmInformerFactory map[string]k8sInformers.SharedInformerFactory,
podSpecPatch *apiv1.PodSpec,
) (executortype.ExecutorType, error) {
@@ -148,11 +147,11 @@ func MakeNewDeploy(
nd.svcLister[ns] = informerFactory.Core().V1().Services().Lister()
nd.svcListerSynced[ns] = informerFactory.Core().V1().Services().Informer().HasSynced
}
for _, fnInformer := range funcInformer {
fnInformer.Informer().AddEventHandler(nd.FunctionEventHandlers(ctx))
for _, factory := range finformerFactory {
factory.Core().V1().Functions().Informer().AddEventHandler(nd.FunctionEventHandlers(ctx))
}
for _, envInformer := range envInformer {
envInformer.Informer().AddEventHandler(nd.EnvEventHandlers(ctx))
for _, factory := range finformerFactory {
factory.Core().V1().Environments().Informer().AddEventHandler(nd.EnvEventHandlers(ctx))
}
return nd, nil
}
@@ -21,7 +21,6 @@ import (
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"
finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/loggerfactory"
)
@@ -35,25 +34,13 @@ const (
configmapName string = "newdeploy-test-configmap"
)
func runInformers(ctx context.Context, informers []k8sCache.SharedIndexInformer) {
// Run all informers
for _, informer := range informers {
go informer.Run(ctx.Done())
}
}
func TestRefreshFuncPods(t *testing.T) {
os.Setenv("DEBUG_ENV", "true")
logger := loggerfactory.GetLogger()
kubernetesClient := fake.NewSimpleClientset()
fissionClient := fClient.NewSimpleClientset()
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30)
funcInformer := map[string]finformerv1.FunctionInformer{
metav1.NamespaceAll: informerFactory.Core().V1().Functions(),
}
envInformer := map[string]finformerv1.EnvironmentInformer{
metav1.NamespaceAll: informerFactory.Core().V1().Environments(),
}
factory := make(map[string]genInformer.SharedInformerFactory, 0)
factory[metav1.NamespaceAll] = genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30)
executorLabel, err := utils.GetInformerLabelByExecutor(fv1.ExecutorTypeNewdeploy)
if err != nil {
@@ -70,7 +57,7 @@ func TestRefreshFuncPods(t *testing.T) {
}
executor, err := MakeNewDeploy(ctx, logger, fissionClient, kubernetesClient, fetcherConfig, "test",
funcInformer, envInformer, ndmInformerFactory, nil)
factory, ndmInformerFactory, nil)
if err != nil {
t.Fatalf("new deploy manager creation failed: %s", err)
}
@@ -87,16 +74,13 @@ func TestRefreshFuncPods(t *testing.T) {
go ndm.Run(ctx)
t.Log("New deploy manager started")
informer := []k8sCache.SharedIndexInformer{
envInformer[metav1.NamespaceAll].Informer(),
funcInformer[metav1.NamespaceAll].Informer(),
for _, f := range factory {
f.Start(ctx.Done())
}
for _, informerFactory := range ndmInformerFactory {
informer = append(informer, informerFactory.Apps().V1().Deployments().Informer())
informer = append(informer, informerFactory.Core().V1().Services().Informer())
informerFactory.Start(ctx.Done())
}
runInformers(ctx, informer)
t.Log("Informers required for new deploy manager started")
waitSynced := make([]k8sCache.InformerSynced, 0)
@@ -59,12 +59,6 @@ func FunctionEventHandlers(ctx context.Context, logger *zap.Logger, kubernetesCl
return
}
// create or update role-binding
envNs := fissionfnNamespace
if fn.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = fn.Spec.Environment.Namespace
}
if istioEnabled {
// create a same name service for function
// since istio only allows the traffic to service
@@ -74,6 +68,7 @@ func FunctionEventHandlers(ctx context.Context, logger *zap.Logger, kubernetesCl
}
svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace)
envNs := utils.DefaultNSResolver().GetFunctionNS(fn.Spec.Environment.Namespace)
// service for accepting user traffic
svc := apiv1.Service{
@@ -124,12 +119,9 @@ func FunctionEventHandlers(ctx context.Context, logger *zap.Logger, kubernetesCl
return
}
envNs := fissionfnNamespace
if fn.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = fn.Spec.Environment.Namespace
}
if istioEnabled {
envNs := utils.DefaultNSResolver().GetFunctionNS(fn.Spec.Environment.Namespace)
svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace)
// delete function istio service
err := kubernetesClient.CoreV1().Services(envNs).Delete(ctx, svcName, metav1.DeleteOptions{})
+3 -4
View File
@@ -50,7 +50,7 @@ import (
executorUtils "github.com/fission/fission/pkg/executor/util"
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"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
"github.com/fission/fission/pkg/utils"
otelUtils "github.com/fission/fission/pkg/utils/otel"
)
@@ -117,8 +117,7 @@ func MakeGenericPoolManager(ctx context.Context,
metricsClient metricsclient.Interface,
fetcherConfig *fetcherConfig.Config,
instanceID string,
funcInformer map[string]finformerv1.FunctionInformer,
envInformer map[string]finformerv1.EnvironmentInformer,
finformerFactory map[string]genInformer.SharedInformerFactory,
gpmInformerFactory map[string]k8sInformers.SharedInformerFactory,
podSpecPatch *apiv1.PodSpec,
) (executortype.ExecutorType, error) {
@@ -135,7 +134,7 @@ func MakeGenericPoolManager(ctx context.Context,
}
poolPodC := NewPoolPodController(ctx, gpmLogger, kubernetesClient,
enableIstio, funcInformer, envInformer, gpmInformerFactory)
enableIstio, finformerFactory, gpmInformerFactory)
gpm := &GenericPoolManager{
logger: gpmLogger,
@@ -37,7 +37,7 @@ import (
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"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
flisterv1 "github.com/fission/fission/pkg/generated/listers/core/v1"
"github.com/fission/fission/pkg/utils"
)
@@ -70,8 +70,7 @@ type (
func NewPoolPodController(ctx context.Context, logger *zap.Logger,
kubernetesClient kubernetes.Interface,
enableIstio bool,
funcInformer map[string]finformerv1.FunctionInformer,
envInformer map[string]finformerv1.EnvironmentInformer,
finformerFactory map[string]genInformer.SharedInformerFactory,
gpmInformerFactory map[string]k8sInformers.SharedInformerFactory) *PoolPodController {
logger = logger.Named("pool_pod_controller")
p := &PoolPodController{
@@ -87,17 +86,19 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger,
envDeleteQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "EnvDeleteQueue"),
spCleanupPodQueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "SpecializedPodCleanupQueue"),
}
for _, informer := range funcInformer {
informer.Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace), p.enableIstio))
if p.enableIstio {
for _, factory := range finformerFactory {
factory.Core().V1().Functions().Informer().AddEventHandler(FunctionEventHandlers(ctx, p.logger, p.kubernetesClient, p.nsResolver.ResolveNamespace(p.nsResolver.FunctionNamespace), p.enableIstio))
}
}
for ns, informer := range envInformer {
informer.Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
for ns, informer := range finformerFactory {
informer.Core().V1().Environments().Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
AddFunc: p.enqueueEnvAdd,
UpdateFunc: p.enqueueEnvUpdate,
DeleteFunc: p.enqueueEnvDelete,
})
p.envLister[ns] = informer.Lister()
p.envListerSynced[ns] = informer.Informer().HasSynced
p.envLister[ns] = informer.Core().V1().Environments().Lister()
p.envListerSynced[ns] = informer.Core().V1().Environments().Informer().HasSynced
}
for ns, informerFactory := range gpmInformerFactory {
informerFactory.Apps().V1().ReplicaSets().Informer().AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
@@ -105,8 +106,8 @@ func NewPoolPodController(ctx context.Context, logger *zap.Logger,
UpdateFunc: p.handleRSUpdate,
DeleteFunc: p.handleRSDelete,
})
p.podLister[ns] = informerFactory.Core().V1().Pods().Lister()
p.podListerSynced[ns] = informerFactory.Core().V1().Pods().Informer().HasSynced
p.podLister[ns] = informerFactory.Core().V1().Pods().Lister()
}
p.logger.Info("pool pod controller handlers registered")
@@ -32,34 +32,18 @@ import (
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"
finformerv1 "github.com/fission/fission/pkg/generated/informers/externalversions/core/v1"
"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) {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
logger := loggerfactory.GetLogger()
kubernetesClient := fake.NewSimpleClientset()
fissionClient := fClient.NewSimpleClientset()
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Minute*30)
funcInformer := map[string]finformerv1.FunctionInformer{
metav1.NamespaceAll: informerFactory.Core().V1().Functions(),
}
pkgInformer := map[string]finformerv1.PackageInformer{
metav1.NamespaceAll: informerFactory.Core().V1().Packages(),
}
envInformer := map[string]finformerv1.EnvironmentInformer{
metav1.NamespaceAll: informerFactory.Core().V1().Environments(),
}
factory := make(map[string]genInformer.SharedInformerFactory, 0)
factory[metav1.NamespaceDefault] = genInformer.NewFilteredSharedInformerFactory(fissionClient, time.Minute*30, metav1.NamespaceDefault, nil)
executorLabel, err := utils.GetInformerLabelByExecutor(fv1.ExecutorTypePoolmgr)
if err != nil {
@@ -68,9 +52,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
gpmInformerFactory := utils.GetInformerFactoryByExecutor(kubernetesClient, executorLabel, time.Minute*30)
ppc := NewPoolPodController(ctx, logger, kubernetesClient, false,
funcInformer,
envInformer,
gpmInformerFactory)
factory, gpmInformerFactory)
executorInstanceID := strings.ToLower(uniuri.NewLen(8))
metricsClient := metricsclient.NewSimpleClientset()
@@ -82,8 +64,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
logger,
fissionClient, kubernetesClient, metricsClient,
fetcherConfig, executorInstanceID,
funcInformer, envInformer,
gpmInformerFactory, nil)
factory, gpmInformerFactory, nil)
if err != nil {
t.Fatalf("Error creating generic pool manager: %v", err)
}
@@ -92,23 +73,18 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
go ppc.Run(ctx, ctx.Done())
informers := []k8sCache.SharedIndexInformer{
funcInformer[metav1.NamespaceAll].Informer(),
pkgInformer[metav1.NamespaceAll].Informer(),
envInformer[metav1.NamespaceAll].Informer(),
for _, f := range factory {
f.Start(ctx.Done())
}
for _, informerFactory := range gpmInformerFactory {
informers = append(informers, informerFactory.Core().V1().Pods().Informer())
informers = append(informers, informerFactory.Apps().V1().ReplicaSets().Informer())
informerFactory.Start(ctx.Done())
}
runInformers(ctx, informers)
pod := &corev1.Pod{
ObjectMeta: metav1.ObjectMeta{
Name: "test-pod",
Namespace: "test-different-namespace",
Namespace: metav1.NamespaceDefault,
},
Status: corev1.PodStatus{
Phase: corev1.PodRunning,
+2 -2
View File
@@ -217,7 +217,7 @@ func getEnvVarlist(ctx context.Context, mqt *fv1.MessageQueueTrigger, routerURL
// Add Auth Fields
secretName := mqt.Spec.Secret
if len(secretName) > 0 {
secret, err := kubeClient.CoreV1().Secrets(apiv1.NamespaceDefault).Get(ctx, secretName, metav1.GetOptions{})
secret, err := kubeClient.CoreV1().Secrets(mqt.Namespace).Get(ctx, secretName, metav1.GetOptions{})
if err != nil {
return nil, err
}
@@ -304,7 +304,7 @@ func checkAndUpdateTriggerFields(mqt, newMqt *fv1.MessageQueueTrigger) bool {
}
func getAuthTriggerSpec(ctx context.Context, mqt *fv1.MessageQueueTrigger, authenticationRef string, kubeClient kubernetes.Interface) (*unstructured.Unstructured, error) {
secret, err := kubeClient.CoreV1().Secrets(apiv1.NamespaceDefault).Get(ctx, mqt.Spec.Secret, metav1.GetOptions{})
secret, err := kubeClient.CoreV1().Secrets(mqt.Namespace).Get(ctx, mqt.Spec.Secret, metav1.GetOptions{})
if err != nil {
return nil, err
}