Shared informers (#2092)

* Use shared informers across executor for fission CRs

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>

* Add informer for configmap and secrets

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>

* Use shared informers in router

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>

* Add shared informer for canary config

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>

* Add sharedinformer for ready pod check

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>

* Refer informer instead of store in resolver

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>

* Run informers before executors

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>

* Add missing canary config handler calls

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>

* Enable websocket test

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>

* Run dump collection always

Signed-off-by: Sanket Sudake <sanketsudake@gmail.com>
This commit is contained in:
Sanket Sudake
2021-07-05 20:28:26 +05:30
committed by GitHub
parent 43814450e5
commit f86fde81e6
22 changed files with 793 additions and 718 deletions
+1
View File
@@ -108,6 +108,7 @@ jobs:
run: ./test/kind_CI.sh
- name: Collect Fission Dump
if: ${{ always() }}
run: |
command -v fission && fission support dump
+1
View File
@@ -84,6 +84,7 @@ jobs:
source ./test/upgrade_test/fission_objects.sh test_fission_objects
- name: Collect Fission Dump
if: ${{ always() }}
run: |
command -v fission && fission support dump
+34 -40
View File
@@ -28,13 +28,13 @@ import (
"go.uber.org/zap"
k8serrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/fields"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
)
const (
@@ -45,8 +45,7 @@ type canaryConfigMgr struct {
logger *zap.Logger
fissionClient *crd.FissionClient
kubeClient *kubernetes.Clientset
canaryConfigStore k8sCache.Store
canaryConfigController k8sCache.Controller
canaryConfigInformer *k8sCache.SharedIndexInformer
promClient *PrometheusApiClient
crdClient rest.Interface
canaryCfgCancelFuncMap *canaryConfigCancelFuncMap
@@ -97,49 +96,44 @@ func MakeCanaryConfigMgr(logger *zap.Logger, fissionClient *crd.FissionClient, k
canaryCfgCancelFuncMap: makecanaryConfigCancelFuncMap(),
}
store, controller := configMgr.initCanaryConfigController()
configMgr.canaryConfigStore = store
configMgr.canaryConfigController = controller
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Second*30)
informer := informerFactory.Core().V1().CanaryConfigs().Informer()
configMgr.canaryConfigInformer = &informer
configMgr.CanaryConfigEventHandlers()
return configMgr, nil
}
func (canaryCfgMgr *canaryConfigMgr) initCanaryConfigController() (k8sCache.Store, k8sCache.Controller) {
resyncPeriod := 30 * time.Second
listWatch := k8sCache.NewListWatchFromClient(canaryCfgMgr.crdClient, "canaryconfigs", metav1.NamespaceAll, fields.Everything())
store, controller := k8sCache.NewInformer(listWatch, &fv1.CanaryConfig{}, resyncPeriod,
k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
canaryConfig := obj.(*fv1.CanaryConfig)
if canaryConfig.Status.Status == fv1.CanaryConfigStatusPending {
go canaryCfgMgr.addCanaryConfig(canaryConfig)
}
},
DeleteFunc: func(obj interface{}) {
canaryConfig := obj.(*fv1.CanaryConfig)
go canaryCfgMgr.deleteCanaryConfig(canaryConfig)
},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
oldConfig := oldObj.(*fv1.CanaryConfig)
newConfig := newObj.(*fv1.CanaryConfig)
if oldConfig.ObjectMeta.ResourceVersion != newConfig.ObjectMeta.ResourceVersion &&
newConfig.Status.Status == fv1.CanaryConfigStatusPending {
canaryCfgMgr.logger.Info("update canary config invoked",
zap.String("name", newConfig.ObjectMeta.Name),
zap.String("namespace", newConfig.ObjectMeta.Namespace),
zap.String("version", newConfig.ObjectMeta.ResourceVersion))
go canaryCfgMgr.updateCanaryConfig(oldConfig, newConfig)
}
go canaryCfgMgr.reSyncCanaryConfigs()
func (canaryCfgMgr *canaryConfigMgr) CanaryConfigEventHandlers() {
(*canaryCfgMgr.canaryConfigInformer).AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
canaryConfig := obj.(*fv1.CanaryConfig)
if canaryConfig.Status.Status == fv1.CanaryConfigStatusPending {
go canaryCfgMgr.addCanaryConfig(canaryConfig)
}
},
DeleteFunc: func(obj interface{}) {
canaryConfig := obj.(*fv1.CanaryConfig)
go canaryCfgMgr.deleteCanaryConfig(canaryConfig)
},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
oldConfig := oldObj.(*fv1.CanaryConfig)
newConfig := newObj.(*fv1.CanaryConfig)
if oldConfig.ObjectMeta.ResourceVersion != newConfig.ObjectMeta.ResourceVersion &&
newConfig.Status.Status == fv1.CanaryConfigStatusPending {
canaryCfgMgr.logger.Info("update canary config invoked",
zap.String("name", newConfig.ObjectMeta.Name),
zap.String("namespace", newConfig.ObjectMeta.Namespace),
zap.String("version", newConfig.ObjectMeta.ResourceVersion))
go canaryCfgMgr.updateCanaryConfig(oldConfig, newConfig)
}
go canaryCfgMgr.reSyncCanaryConfigs()
},
})
return store, controller
},
})
}
func (canaryCfgMgr *canaryConfigMgr) Run(ctx context.Context) {
go canaryCfgMgr.canaryConfigController.Run(ctx.Done())
go (*canaryCfgMgr.canaryConfigInformer).Run(ctx.Done())
canaryCfgMgr.logger.Info("started canary configmgr controller")
}
@@ -501,7 +495,7 @@ func (canaryCfgMgr *canaryConfigMgr) rollForward(canaryConfig *fv1.CanaryConfig,
}
func (canaryCfgMgr *canaryConfigMgr) reSyncCanaryConfigs() {
for _, obj := range canaryCfgMgr.canaryConfigStore.List() {
for _, obj := range (*canaryCfgMgr.canaryConfigInformer).GetStore().List() {
canaryConfig := obj.(*fv1.CanaryConfig)
_, err := canaryCfgMgr.canaryCfgCancelFuncMap.lookup(&canaryConfig.ObjectMeta)
if err != nil && canaryConfig.Status.Status == fv1.CanaryConfigStatusPending {
+73
View File
@@ -0,0 +1,73 @@
/*
Copyright 2021 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 cms
import (
"context"
"go.uber.org/zap"
apiv1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/executor/executortype"
)
func getConfigmapRelatedFuncs(logger *zap.Logger, m *metav1.ObjectMeta, fissionClient *crd.FissionClient) ([]fv1.Function, error) {
funcList, err := fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{})
if err != nil {
return nil, err
}
// In future a cache that populates at start and is updated on changes might be better solution
relatedFunctions := make([]fv1.Function, 0)
for _, f := range funcList.Items {
for _, cm := range f.Spec.ConfigMaps {
if (cm.Name == m.Name) && (cm.Namespace == m.Namespace) {
relatedFunctions = append(relatedFunctions, f)
break
}
}
}
return relatedFunctions, nil
}
func ConfigMapEventHandlers(logger *zap.Logger, fissionClient *crd.FissionClient,
kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType) k8sCache.ResourceEventHandlerFuncs {
return k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {},
DeleteFunc: func(obj interface{}) {},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
oldCm := oldObj.(*apiv1.ConfigMap)
newCm := newObj.(*apiv1.ConfigMap)
if oldCm.ObjectMeta.ResourceVersion != newCm.ObjectMeta.ResourceVersion {
if newCm.ObjectMeta.Namespace != "kube-system" {
logger.Debug("Configmap changed",
zap.String("configmap_name", newCm.ObjectMeta.Name),
zap.String("configmap_namespace", newCm.ObjectMeta.Namespace))
}
funcs, err := getConfigmapRelatedFuncs(logger, &newCm.ObjectMeta, fissionClient)
if err != nil {
logger.Error("Failed to get functions related to configmap", zap.String("configmap_name", newCm.ObjectMeta.Name), zap.String("configmap_namespace", newCm.ObjectMeta.Namespace))
}
refreshPods(logger, funcs, types)
}
},
}
}
+14 -115
View File
@@ -17,16 +17,10 @@ limitations under the License.
package cms
import (
"context"
"time"
"github.com/pkg/errors"
"go.uber.org/zap"
apiv1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/fields"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/tools/cache"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd"
@@ -38,126 +32,31 @@ type (
ConfigSecretController struct {
logger *zap.Logger
configmapController cache.Controller
secretController cache.Controller
configmapInformer *k8sCache.SharedIndexInformer
secretInformer *k8sCache.SharedIndexInformer
fissionClient *crd.FissionClient
}
)
//MakeConfigSecretController makes a controller for configmaps and secrets which changes related functions
// MakeConfigSecretController makes a controller for configmaps and secrets which changes related functions
func MakeConfigSecretController(logger *zap.Logger, fissionClient *crd.FissionClient,
kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType) *ConfigSecretController {
kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType,
configmapInformer *k8sCache.SharedIndexInformer,
secretInformer *k8sCache.SharedIndexInformer) *ConfigSecretController {
logger.Debug("Creating ConfigMap & Secret Controller")
_, cmcontroller := initConfigmapController(logger, fissionClient, kubernetesClient, types)
_, scontroller := initSecretController(logger, fissionClient, kubernetesClient, types)
cmsController := &ConfigSecretController{
logger: logger,
configmapController: cmcontroller,
secretController: scontroller,
fissionClient: fissionClient,
logger: logger,
configmapInformer: configmapInformer,
secretInformer: secretInformer,
fissionClient: fissionClient,
}
(*configmapInformer).AddEventHandler(ConfigMapEventHandlers(logger, fissionClient, kubernetesClient, types))
(*secretInformer).AddEventHandler(SecretEventHandlers(logger, fissionClient, kubernetesClient, types))
return cmsController
}
//Run runs the controllers for configmaps and secrets
func (csController *ConfigSecretController) Run(ctx context.Context) {
go csController.configmapController.Run(ctx.Done())
go csController.secretController.Run(ctx.Done())
}
func initConfigmapController(logger *zap.Logger, fissionClient *crd.FissionClient,
kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType) (cache.Store, cache.Controller) {
resyncPeriod := 30 * time.Second
listWatch := cache.NewListWatchFromClient(kubernetesClient.CoreV1().RESTClient(), "configmaps", metav1.NamespaceAll, fields.Everything())
store, controller := cache.NewInformer(listWatch, &apiv1.ConfigMap{}, resyncPeriod, cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {},
DeleteFunc: func(obj interface{}) {},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
oldCm := oldObj.(*apiv1.ConfigMap)
newCm := newObj.(*apiv1.ConfigMap)
if oldCm.ObjectMeta.ResourceVersion != newCm.ObjectMeta.ResourceVersion {
if newCm.ObjectMeta.Namespace != "kube-system" {
logger.Debug("Configmap changed",
zap.String("configmap_name", newCm.ObjectMeta.Name),
zap.String("configmap_namespace", newCm.ObjectMeta.Namespace))
}
funcs, err := getConfigmapRelatedFuncs(logger, &newCm.ObjectMeta, fissionClient)
if err != nil {
logger.Error("Failed to get functions related to configmap", zap.String("configmap_name", newCm.ObjectMeta.Name), zap.String("configmap_namespace", newCm.ObjectMeta.Namespace))
}
refreshPods(logger, funcs, types)
}
},
})
return store, controller
}
func getConfigmapRelatedFuncs(logger *zap.Logger, m *metav1.ObjectMeta, fissionClient *crd.FissionClient) ([]fv1.Function, error) {
funcList, err := fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{})
if err != nil {
return nil, err
}
// In future a cache that populates at start and is updated on changes might be better solution
relatedFunctions := make([]fv1.Function, 0)
for _, f := range funcList.Items {
for _, cm := range f.Spec.ConfigMaps {
if (cm.Name == m.Name) && (cm.Namespace == m.Namespace) {
relatedFunctions = append(relatedFunctions, f)
break
}
}
}
return relatedFunctions, nil
}
func initSecretController(logger *zap.Logger, fissionClient *crd.FissionClient,
kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType) (cache.Store, cache.Controller) {
resyncPeriod := 30 * time.Second
listWatch := cache.NewListWatchFromClient(kubernetesClient.CoreV1().RESTClient(), "secrets", metav1.NamespaceAll, fields.Everything())
store, controller := cache.NewInformer(listWatch, &apiv1.Secret{}, resyncPeriod, cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {},
DeleteFunc: func(obj interface{}) {},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
oldS := oldObj.(*apiv1.Secret)
newS := newObj.(*apiv1.Secret)
if oldS.ObjectMeta.ResourceVersion != newS.ObjectMeta.ResourceVersion {
if newS.ObjectMeta.Namespace != "kube-system" {
logger.Debug("Secret changed",
zap.String("configmap_name", newS.ObjectMeta.Name),
zap.String("configmap_namespace", newS.ObjectMeta.Namespace))
}
funcs, err := getSecretRelatedFuncs(logger, &newS.ObjectMeta, fissionClient)
if err != nil {
logger.Error("Failed to get functions related to secret", zap.String("secret_name", newS.ObjectMeta.Name), zap.String("secret_namespace", newS.ObjectMeta.Namespace))
}
refreshPods(logger, funcs, types)
}
},
})
return store, controller
}
func getSecretRelatedFuncs(logger *zap.Logger, m *metav1.ObjectMeta, fissionClient *crd.FissionClient) ([]fv1.Function, error) {
funcList, err := fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{})
if err != nil {
return nil, err
}
// In future a cache that populates at start and is updated on changes might be better solution
relatedFunctions := make([]fv1.Function, 0)
for _, f := range funcList.Items {
for _, secret := range f.Spec.Secrets {
if (secret.Name == m.Name) && (secret.Namespace == m.Namespace) {
relatedFunctions = append(relatedFunctions, f)
break
}
}
}
return relatedFunctions, nil
}
func refreshPods(logger *zap.Logger, funcs []fv1.Function, types map[fv1.ExecutorType]executortype.ExecutorType) {
for _, f := range funcs {
var err error
+72
View File
@@ -0,0 +1,72 @@
/*
Copyright 2021 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 cms
import (
"context"
"go.uber.org/zap"
apiv1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/executor/executortype"
)
func getSecretRelatedFuncs(logger *zap.Logger, m *metav1.ObjectMeta, fissionClient *crd.FissionClient) ([]fv1.Function, error) {
funcList, err := fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{})
if err != nil {
return nil, err
}
// In future a cache that populates at start and is updated on changes might be better solution
relatedFunctions := make([]fv1.Function, 0)
for _, f := range funcList.Items {
for _, secret := range f.Spec.Secrets {
if (secret.Name == m.Name) && (secret.Namespace == m.Namespace) {
relatedFunctions = append(relatedFunctions, f)
break
}
}
}
return relatedFunctions, nil
}
func SecretEventHandlers(logger *zap.Logger, fissionClient *crd.FissionClient,
kubernetesClient *kubernetes.Clientset, types map[fv1.ExecutorType]executortype.ExecutorType) k8sCache.ResourceEventHandlerFuncs {
return k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {},
DeleteFunc: func(obj interface{}) {},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
oldS := oldObj.(*apiv1.Secret)
newS := newObj.(*apiv1.Secret)
if oldS.ObjectMeta.ResourceVersion != newS.ObjectMeta.ResourceVersion {
if newS.ObjectMeta.Namespace != "kube-system" {
logger.Debug("Secret changed",
zap.String("configmap_name", newS.ObjectMeta.Name),
zap.String("configmap_namespace", newS.ObjectMeta.Namespace))
}
funcs, err := getSecretRelatedFuncs(logger, &newS.ObjectMeta, fissionClient)
if err != nil {
logger.Error("Failed to get functions related to secret", zap.String("secret_name", newS.ObjectMeta.Name), zap.String("secret_namespace", newS.ObjectMeta.Namespace))
}
refreshPods(logger, funcs, types)
}
},
}
}
+34 -8
View File
@@ -30,6 +30,8 @@ import (
"github.com/pkg/errors"
"github.com/prometheus/client_golang/prometheus/promhttp"
"go.uber.org/zap"
k8sInformers "k8s.io/client-go/informers"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd"
@@ -41,6 +43,7 @@ import (
"github.com/fission/fission/pkg/executor/reaper"
"github.com/fission/fission/pkg/executor/util"
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
)
type (
@@ -69,8 +72,9 @@ type (
)
// MakeExecutor returns an Executor for given ExecutorType(s).
func MakeExecutor(logger *zap.Logger, cms *cms.ConfigSecretController,
fissionClient *crd.FissionClient, types map[fv1.ExecutorType]executortype.ExecutorType) (*Executor, error) {
func MakeExecutor(ctx context.Context, logger *zap.Logger, cms *cms.ConfigSecretController,
fissionClient *crd.FissionClient, types map[fv1.ExecutorType]executortype.ExecutorType,
informers []k8sCache.SharedIndexInformer) (*Executor, error) {
executor := &Executor{
logger: logger.Named("executor"),
cms: cms,
@@ -80,12 +84,18 @@ func MakeExecutor(logger *zap.Logger, cms *cms.ConfigSecretController,
requestChan: make(chan *createFuncServiceRequest),
fsCreateWg: make(map[string]*sync.WaitGroup),
}
// Run all informers
for _, informer := range informers {
go informer.Run(ctx.Done())
}
for _, et := range types {
go func(et executortype.ExecutorType) {
et.Run(context.Background())
et.Run(ctx)
}(et)
}
go cms.Run(context.Background())
go executor.serveCreateFuncServices()
return executor, nil
@@ -257,10 +267,17 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames
logger.Info("Starting executor", zap.String("instanceID", executorInstanceID))
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Second*30)
funcInformer := informerFactory.Core().V1().Functions().Informer()
pkgInformer := informerFactory.Core().V1().Packages().Informer()
envInformer := informerFactory.Core().V1().Environments().Informer()
gpm, err := poolmgr.MakeGenericPoolManager(
logger,
fissionClient, kubernetesClient, metricsClient,
functionNamespace, fetcherConfig, executorInstanceID)
functionNamespace, fetcherConfig, executorInstanceID,
&funcInformer, &pkgInformer,
)
if err != nil {
return errors.Wrap(err, "pool manager creation faied")
}
@@ -268,7 +285,9 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames
ndm, err := newdeploy.MakeNewDeploy(
logger,
fissionClient, kubernetesClient, fissionClient.CoreV1().RESTClient(),
functionNamespace, fetcherConfig, executorInstanceID)
functionNamespace, fetcherConfig, executorInstanceID,
&funcInformer, &envInformer,
)
if err != nil {
return errors.Wrap(err, "new deploy manager creation faied")
}
@@ -294,9 +313,16 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames
// TODO: use context to control the waiting time once kubernetes client supports it.
util.WaitTimeout(wg, 30*time.Second)
cms := cms.MakeConfigSecretController(logger, fissionClient, kubernetesClient, executorTypes)
k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Second*30)
configmapInformer := k8sInformerFactory.Core().V1().ConfigMaps().Informer()
secretInformer := k8sInformerFactory.Core().V1().Secrets().Informer()
api, err := MakeExecutor(logger, cms, fissionClient, executorTypes)
cms := cms.MakeConfigSecretController(logger, fissionClient, kubernetesClient, executorTypes, &configmapInformer, &secretInformer)
ctx := context.Background()
api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes, []k8sCache.SharedIndexInformer{
funcInformer, pkgInformer, envInformer, configmapInformer, secretInformer,
})
if err != nil {
return err
}
@@ -0,0 +1,54 @@
/*
Copyright 2021 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 newdeploy
import (
"context"
"go.uber.org/zap"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
)
func (deploy *NewDeploy) EnvEventHandlers() k8sCache.ResourceEventHandlerFuncs {
return k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {},
DeleteFunc: func(obj interface{}) {},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
newEnv := newObj.(*fv1.Environment)
oldEnv := oldObj.(*fv1.Environment)
// Currently only an image update in environment calls for function's deployment recreation. In future there might be more attributes which would want to do it
if oldEnv.Spec.Runtime.Image != newEnv.Spec.Runtime.Image {
deploy.logger.Debug("Updating all function of the environment that changed, old env:", zap.Any("environment", oldEnv))
funcs := deploy.getEnvFunctions(&newEnv.ObjectMeta)
for _, f := range funcs {
function, err := deploy.fissionClient.CoreV1().Functions(f.ObjectMeta.Namespace).Get(context.TODO(), f.ObjectMeta.Name, metav1.GetOptions{})
if err != nil {
deploy.logger.Error("Error getting function", zap.Error(err), zap.Any("function", function))
continue
}
err = deploy.updateFuncDeployment(function, newEnv)
if err != nil {
deploy.logger.Error("Error updating function", zap.Error(err), zap.Any("function", function))
continue
}
}
}
},
}
}
@@ -0,0 +1,68 @@
/*
Copyright 2021 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 newdeploy
import (
"go.uber.org/zap"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
)
func (deploy *NewDeploy) FunctionEventHandlers() k8sCache.ResourceEventHandlerFuncs {
return k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
// TODO: A workaround to process items in parallel. We should use workqueue ("k8s.io/client-go/util/workqueue")
// and worker pattern to process items instead of moving process to another goroutine.
// example: https://github.com/kubernetes/kubernetes/blob/master/pkg/controller/job/job_controller.go
go func() {
fn := obj.(*fv1.Function)
deploy.logger.Debug("create deployment for function", zap.Any("fn", fn.ObjectMeta), zap.Any("fnspec", fn.Spec))
_, err := deploy.createFunction(fn)
if err != nil {
deploy.logger.Error("error eager creating function",
zap.Error(err),
zap.Any("function", fn))
}
deploy.logger.Debug("end create deployment for function", zap.Any("fn", fn.ObjectMeta), zap.Any("fnspec", fn.Spec))
}()
},
DeleteFunc: func(obj interface{}) {
fn := obj.(*fv1.Function)
go func() {
err := deploy.deleteFunction(fn)
if err != nil {
deploy.logger.Error("error deleting function",
zap.Error(err),
zap.Any("function", fn))
}
}()
},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
oldFn := oldObj.(*fv1.Function)
newFn := newObj.(*fv1.Function)
go func() {
err := deploy.updateFunction(oldFn, newFn)
if err != nil {
deploy.logger.Error("error updating function",
zap.Error(err),
zap.Any("old_function", oldFn),
zap.Any("new_function", newFn))
}
}()
},
}
}
@@ -33,7 +33,6 @@ import (
apiv1 "k8s.io/api/core/v1"
k8sErrs "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/fields"
"k8s.io/apimachinery/pkg/labels"
k8sTypes "k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/kubernetes"
@@ -69,12 +68,11 @@ type (
fsCache *fscache.FunctionServiceCache // cache funcSvc's by function, address and pod name
throttler *throttler.Throttler
funcStore k8sCache.Store
funcController k8sCache.Controller
throttler *throttler.Throttler
funcInformer *k8sCache.SharedIndexInformer
envInformer *k8sCache.SharedIndexInformer
envStore k8sCache.Store
envController k8sCache.Controller
serviceInformer k8sCache.SharedIndexInformer
deploymentInformer k8sCache.SharedIndexInformer
@@ -91,6 +89,8 @@ func MakeNewDeploy(
namespace string,
fetcherConfig *fetcherConfig.Config,
instanceID string,
funcInformer *k8sCache.SharedIndexInformer,
envInformer *k8sCache.SharedIndexInformer,
) (executortype.ExecutorType, error) {
enableIstio := false
if len(os.Getenv("ENABLE_ISTIO")) > 0 {
@@ -118,19 +118,16 @@ func MakeNewDeploy(
useIstio: enableIstio,
defaultIdlePodReapTime: 2 * time.Minute,
funcInformer: funcInformer,
envInformer: envInformer,
}
if nd.crdClient != nil {
fnStore, fnController := nd.initFuncController()
nd.funcStore = fnStore
nd.funcController = fnController
envStore, envController := nd.initEnvController()
nd.envStore = envStore
nd.envController = envController
(*nd.funcInformer).AddEventHandler(nd.FunctionEventHandlers())
(*nd.envInformer).AddEventHandler(nd.EnvEventHandlers())
}
informerFactory, err := utils.GetInformerFacoryByExecutor(nd.kubernetesClient, fv1.ExecutorTypePoolmgr)
informerFactory, err := utils.GetInformerFactoryByExecutor(nd.kubernetesClient, fv1.ExecutorTypePoolmgr)
if err != nil {
return nil, err
}
@@ -141,8 +138,6 @@ func MakeNewDeploy(
// Run start the function and environment controller along with an object reaper.
func (deploy *NewDeploy) Run(ctx context.Context) {
go deploy.funcController.Run(ctx.Done())
go deploy.envController.Run(ctx.Done())
go deploy.serviceInformer.Run(ctx.Done())
go deploy.deploymentInformer.Run(ctx.Done())
go deploy.idleObjectReaper()
@@ -191,19 +186,8 @@ func (deploy *NewDeploy) TapService(svcHost string) error {
return nil
}
func getCachedItem(obj apiv1.ObjectReference, informer k8sCache.SharedIndexInformer) (item interface{}, exists bool, err error) {
store := informer.GetStore()
item, exists, err = store.Get(obj)
if err != nil || !exists {
item, exists, err = store.GetByKey(fmt.Sprintf("%s/%s", obj.Namespace, obj.Name))
}
return item, exists, err
}
func (deploy *NewDeploy) getServiceInfo(obj apiv1.ObjectReference) (*apiv1.Service, error) {
item, exists, err := getCachedItem(obj, deploy.serviceInformer)
item, exists, err := utils.GetCachedItem(obj, deploy.serviceInformer)
if err != nil || !exists {
deploy.logger.Debug(
@@ -220,7 +204,7 @@ func (deploy *NewDeploy) getServiceInfo(obj apiv1.ObjectReference) (*apiv1.Servi
}
func (deploy *NewDeploy) getDeploymentInfo(obj apiv1.ObjectReference) (*appsv1.Deployment, error) {
item, exists, err := getCachedItem(obj, deploy.deploymentInformer)
item, exists, err := utils.GetCachedItem(obj, deploy.deploymentInformer)
if err != nil || !exists {
deploy.logger.Debug(
@@ -378,85 +362,6 @@ func (deploy *NewDeploy) CleanupOldExecutorObjects() {
}
}
func (deploy *NewDeploy) initFuncController() (k8sCache.Store, k8sCache.Controller) {
resyncPeriod := 30 * time.Second
listWatch := k8sCache.NewListWatchFromClient(deploy.crdClient, "functions", metav1.NamespaceAll, fields.Everything())
store, controller := k8sCache.NewInformer(listWatch, &fv1.Function{}, resyncPeriod, k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
// TODO: A workaround to process items in parallel. We should use workqueue ("k8s.io/client-go/util/workqueue")
// and worker pattern to process items instead of moving process to another goroutine.
// example: https://github.com/kubernetes/kubernetes/blob/master/pkg/controller/job/job_controller.go
go func() {
fn := obj.(*fv1.Function)
deploy.logger.Debug("create deployment for function", zap.Any("fn", fn.ObjectMeta), zap.Any("fnspec", fn.Spec))
_, err := deploy.createFunction(fn)
if err != nil {
deploy.logger.Error("error eager creating function",
zap.Error(err),
zap.Any("function", fn))
}
deploy.logger.Debug("end create deployment for function", zap.Any("fn", fn.ObjectMeta), zap.Any("fnspec", fn.Spec))
}()
},
DeleteFunc: func(obj interface{}) {
fn := obj.(*fv1.Function)
go func() {
err := deploy.deleteFunction(fn)
if err != nil {
deploy.logger.Error("error deleting function",
zap.Error(err),
zap.Any("function", fn))
}
}()
},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
oldFn := oldObj.(*fv1.Function)
newFn := newObj.(*fv1.Function)
go func() {
err := deploy.updateFunction(oldFn, newFn)
if err != nil {
deploy.logger.Error("error updating function",
zap.Error(err),
zap.Any("old_function", oldFn),
zap.Any("new_function", newFn))
}
}()
},
})
return store, controller
}
func (deploy *NewDeploy) initEnvController() (k8sCache.Store, k8sCache.Controller) {
resyncPeriod := 30 * time.Second
listWatch := k8sCache.NewListWatchFromClient(deploy.crdClient, "environments", metav1.NamespaceAll, fields.Everything())
store, controller := k8sCache.NewInformer(listWatch, &fv1.Environment{}, resyncPeriod, k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {},
DeleteFunc: func(obj interface{}) {},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
newEnv := newObj.(*fv1.Environment)
oldEnv := oldObj.(*fv1.Environment)
// Currently only an image update in environment calls for function's deployment recreation. In future there might be more attributes which would want to do it
if oldEnv.Spec.Runtime.Image != newEnv.Spec.Runtime.Image {
deploy.logger.Debug("Updating all function of the environment that changed, old env:", zap.Any("environment", oldEnv))
funcs := deploy.getEnvFunctions(&newEnv.ObjectMeta)
for _, f := range funcs {
function, err := deploy.fissionClient.CoreV1().Functions(f.ObjectMeta.Namespace).Get(context.TODO(), f.ObjectMeta.Name, metav1.GetOptions{})
if err != nil {
deploy.logger.Error("Error getting function", zap.Error(err), zap.Any("function", function))
continue
}
err = deploy.updateFuncDeployment(function, newEnv)
if err != nil {
deploy.logger.Error("Error updating function", zap.Error(err), zap.Any("function", function))
continue
}
}
}
},
})
return store, controller
}
func (deploy *NewDeploy) getEnvFunctions(m *metav1.ObjectMeta) []fv1.Function {
funcList, err := deploy.fissionClient.CoreV1().Functions(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{})
if err != nil {
@@ -876,7 +781,7 @@ func (deploy *NewDeploy) idleObjectReaper() {
// For function with the environment that no longer exists, executor
// scales down the deployment as usual and prints log to notify user.
if _, ok := envList[fsvc.Environment.ObjectMeta.UID]; !ok {
deploy.logger.Error("function environment no longer exists",
deploy.logger.Warn("function environment no longer exists",
zap.String("environment", fsvc.Environment.ObjectMeta.Name),
zap.String("function", fsvc.Name))
}
@@ -0,0 +1,198 @@
/*
Copyright 2018 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"
"go.uber.org/zap"
apiv1 "k8s.io/api/core/v1"
kerrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/util/intstr"
"k8s.io/client-go/kubernetes"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/utils"
)
func getIstioServiceLabels(fnName string) map[string]string {
return map[string]string{
"functionName": fnName,
}
}
func (gpm *GenericPoolManager) FunctionEventHandlers(kubernetesClient *kubernetes.Clientset, fissionfnNamespace string, istioEnabled bool) k8sCache.ResourceEventHandlerFuncs {
return k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
fn := obj.(*fv1.Function)
// Since istio only allows accessing pod through k8s service,
// for the functions with executor type "poolmgr" we need to
// create a service for sending requests to pod in pool.
// Functions with executor type "Newdeploy" is specialized at
// pod starts. In this case, just ignore such functions.
fnExecutorType := fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType
// In some cases, user may not enter the executorType explicitly, for example in his spec.yaml.
// we assume it to be of type poolmgr
if fnExecutorType != "" && fnExecutorType != fv1.ExecutorTypePoolmgr {
return
}
// create or update role-binding
envNs := fissionfnNamespace
if fn.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = fn.Spec.Environment.Namespace
}
// TODO : Just bring to your attention during review :
// setup rolebinding is tried, if it fails, we don't return. we just log an error and move on, because :
// 1. not all functions have secrets and/or configmaps, so things will work without this rolebinding in that case.
// 2. on the contrary, when the route is tried, the env fetcher logs will show a 403 forbidden message and same will be relayed to executor.
err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs)
if err != nil {
gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB))
} else {
gpm.logger.Debug("successfully set up rolebinding for fetcher service account for function",
zap.String("service_account", fv1.FissionFetcherSA),
zap.String("service_account_namepsace", envNs),
zap.String("function_name", fn.ObjectMeta.Name),
zap.String("function_namespace", fn.ObjectMeta.Namespace))
}
if istioEnabled {
// create a same name service for function
// since istio only allows the traffic to service
sel := map[string]string{
"functionName": fn.ObjectMeta.Name,
"functionUid": string(fn.ObjectMeta.UID),
}
svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace)
// service for accepting user traffic
svc := apiv1.Service{
ObjectMeta: metav1.ObjectMeta{
Namespace: envNs,
Name: svcName,
Labels: getIstioServiceLabels(fn.ObjectMeta.Name),
},
Spec: apiv1.ServiceSpec{
Type: apiv1.ServiceTypeClusterIP,
Ports: []apiv1.ServicePort{
// Service port name should begin with a recognized prefix, or the traffic will be
// treated as TCP traffic. (https://istio.io/docs/setup/kubernetes/additional-setup/requirements/)
{
Name: "http-fetcher",
Protocol: apiv1.ProtocolTCP,
Port: 8000,
TargetPort: intstr.FromInt(8000),
},
{
Name: "http-env",
Protocol: apiv1.ProtocolTCP,
Port: 8888,
TargetPort: intstr.FromInt(8888),
},
},
Selector: sel,
},
}
// create function istio service if it does not exist
_, err = kubernetesClient.CoreV1().Services(envNs).Create(context.TODO(), &svc, metav1.CreateOptions{})
if err != nil && !kerrors.IsAlreadyExists(err) {
gpm.logger.Error("error creating istio service for function",
zap.Error(err),
zap.String("service_name", svcName),
zap.String("function_name", fn.ObjectMeta.Name),
zap.Any("selectors", sel))
}
}
},
DeleteFunc: func(obj interface{}) {
fn := obj.(*fv1.Function)
fnExecutorType := fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType
if fnExecutorType != "" && fnExecutorType != fv1.ExecutorTypePoolmgr {
return
}
envNs := fissionfnNamespace
if fn.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = fn.Spec.Environment.Namespace
}
if istioEnabled {
svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace)
// delete function istio service
err := kubernetesClient.CoreV1().Services(envNs).Delete(context.TODO(), svcName, metav1.DeleteOptions{})
if err != nil && !kerrors.IsNotFound(err) {
gpm.logger.Error("error deleting istio service for function",
zap.Error(err),
zap.String("service_name", svcName),
zap.String("function_name", fn.ObjectMeta.Name))
}
}
},
UpdateFunc: func(oldObj, newObj interface{}) {
oldFunc := oldObj.(*fv1.Function)
newFunc := newObj.(*fv1.Function)
if oldFunc.ObjectMeta.ResourceVersion == newFunc.ObjectMeta.ResourceVersion {
return
}
envChanged := (oldFunc.Spec.Environment.Namespace != newFunc.Spec.Environment.Namespace)
executorTypeChangedToPM := (oldFunc.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType != fv1.ExecutorTypePoolmgr &&
newFunc.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType == fv1.ExecutorTypePoolmgr)
// if a func's env reference gets updated and the newly referenced env is in a different ns,
// we need to create a rolebinding in func's ns so that the fetcher-sa in env ns has access
// to fetch secrets and config maps from the func's ns.
// similarly if executorType changed to Pool Manager, we now need a rolebinding in the func ns for fetcher sa
// present in env ns because for newdeploy, the fetcher sa is in function namespace
if envChanged || executorTypeChangedToPM {
envNs := fissionfnNamespace
if newFunc.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = newFunc.Spec.Environment.Namespace
}
err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.SecretConfigMapGetterRB,
newFunc.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole,
fv1.FissionFetcherSA, envNs)
if err != nil {
gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB))
} else {
gpm.logger.Debug("successfully set up rolebinding for fetcher service account for function",
zap.String("service_account", fv1.FissionFetcherSA),
zap.String("service_account_namepsace", envNs),
zap.String("function_name", newFunc.ObjectMeta.Name),
zap.String("function_namespace", newFunc.ObjectMeta.Namespace))
}
}
},
}
}
@@ -1,208 +0,0 @@
/*
Copyright 2018 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"
"time"
"go.uber.org/zap"
apiv1 "k8s.io/api/core/v1"
kerrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/fields"
"k8s.io/apimachinery/pkg/util/intstr"
"k8s.io/client-go/kubernetes"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/utils"
)
func getIstioServiceLabels(fnName string) map[string]string {
return map[string]string{
"functionName": fnName,
}
}
func (gpm *GenericPoolManager) makeFuncController(fissionClient *crd.FissionClient,
kubernetesClient *kubernetes.Clientset, fissionfnNamespace string, istioEnabled bool) (k8sCache.Store, k8sCache.Controller) {
resyncPeriod := 30 * time.Second
lw := k8sCache.NewListWatchFromClient(fissionClient.CoreV1().RESTClient(), "functions", metav1.NamespaceAll, fields.Everything())
funcStore, controller := k8sCache.NewInformer(lw, &fv1.Function{}, resyncPeriod,
k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
fn := obj.(*fv1.Function)
// Since istio only allows accessing pod through k8s service,
// for the functions with executor type "poolmgr" we need to
// create a service for sending requests to pod in pool.
// Functions with executor type "Newdeploy" is specialized at
// pod starts. In this case, just ignore such functions.
fnExecutorType := fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType
// In some cases, user may not enter the executorType explicitly, for example in his spec.yaml.
// we assume it to be of type poolmgr
if fnExecutorType != "" && fnExecutorType != fv1.ExecutorTypePoolmgr {
return
}
// create or update role-binding
envNs := fissionfnNamespace
if fn.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = fn.Spec.Environment.Namespace
}
// TODO : Just bring to your attention during review :
// setup rolebinding is tried, if it fails, we don't return. we just log an error and move on, because :
// 1. not all functions have secrets and/or configmaps, so things will work without this rolebinding in that case.
// 2. on the contrary, when the route is tried, the env fetcher logs will show a 403 forbidden message and same will be relayed to executor.
err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.SecretConfigMapGetterRB, fn.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs)
if err != nil {
gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB))
} else {
gpm.logger.Debug("successfully set up rolebinding for fetcher service account for function",
zap.String("service_account", fv1.FissionFetcherSA),
zap.String("service_account_namepsace", envNs),
zap.String("function_name", fn.ObjectMeta.Name),
zap.String("function_namespace", fn.ObjectMeta.Namespace))
}
if istioEnabled {
// create a same name service for function
// since istio only allows the traffic to service
sel := map[string]string{
"functionName": fn.ObjectMeta.Name,
"functionUid": string(fn.ObjectMeta.UID),
}
svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace)
// service for accepting user traffic
svc := apiv1.Service{
ObjectMeta: metav1.ObjectMeta{
Namespace: envNs,
Name: svcName,
Labels: getIstioServiceLabels(fn.ObjectMeta.Name),
},
Spec: apiv1.ServiceSpec{
Type: apiv1.ServiceTypeClusterIP,
Ports: []apiv1.ServicePort{
// Service port name should begin with a recognized prefix, or the traffic will be
// treated as TCP traffic. (https://istio.io/docs/setup/kubernetes/additional-setup/requirements/)
{
Name: "http-fetcher",
Protocol: apiv1.ProtocolTCP,
Port: 8000,
TargetPort: intstr.FromInt(8000),
},
{
Name: "http-env",
Protocol: apiv1.ProtocolTCP,
Port: 8888,
TargetPort: intstr.FromInt(8888),
},
},
Selector: sel,
},
}
// create function istio service if it does not exist
_, err = kubernetesClient.CoreV1().Services(envNs).Create(context.TODO(), &svc, metav1.CreateOptions{})
if err != nil && !kerrors.IsAlreadyExists(err) {
gpm.logger.Error("error creating istio service for function",
zap.Error(err),
zap.String("service_name", svcName),
zap.String("function_name", fn.ObjectMeta.Name),
zap.Any("selectors", sel))
}
}
},
DeleteFunc: func(obj interface{}) {
fn := obj.(*fv1.Function)
fnExecutorType := fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType
if fnExecutorType != "" && fnExecutorType != fv1.ExecutorTypePoolmgr {
return
}
envNs := fissionfnNamespace
if fn.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = fn.Spec.Environment.Namespace
}
if istioEnabled {
svcName := utils.GetFunctionIstioServiceName(fn.ObjectMeta.Name, fn.ObjectMeta.Namespace)
// delete function istio service
err := kubernetesClient.CoreV1().Services(envNs).Delete(context.TODO(), svcName, metav1.DeleteOptions{})
if err != nil && !kerrors.IsNotFound(err) {
gpm.logger.Error("error deleting istio service for function",
zap.Error(err),
zap.String("service_name", svcName),
zap.String("function_name", fn.ObjectMeta.Name))
}
}
},
UpdateFunc: func(oldObj, newObj interface{}) {
oldFunc := oldObj.(*fv1.Function)
newFunc := newObj.(*fv1.Function)
if oldFunc.ObjectMeta.ResourceVersion == newFunc.ObjectMeta.ResourceVersion {
return
}
envChanged := (oldFunc.Spec.Environment.Namespace != newFunc.Spec.Environment.Namespace)
executorTypeChangedToPM := (oldFunc.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType != fv1.ExecutorTypePoolmgr &&
newFunc.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType == fv1.ExecutorTypePoolmgr)
// if a func's env reference gets updated and the newly referenced env is in a different ns,
// we need to create a rolebinding in func's ns so that the fetcher-sa in env ns has access
// to fetch secrets and config maps from the func's ns.
// similarly if executorType changed to Pool Manager, we now need a rolebinding in the func ns for fetcher sa
// present in env ns because for newdeploy, the fetcher sa is in function namespace
if envChanged || executorTypeChangedToPM {
envNs := fissionfnNamespace
if newFunc.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = newFunc.Spec.Environment.Namespace
}
err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.SecretConfigMapGetterRB,
newFunc.ObjectMeta.Namespace, fv1.SecretConfigMapGetterCR, fv1.ClusterRole,
fv1.FissionFetcherSA, envNs)
if err != nil {
gpm.logger.Error("error creating rolebinding", zap.Error(err), zap.String("role_binding", fv1.SecretConfigMapGetterRB))
} else {
gpm.logger.Debug("successfully set up rolebinding for fetcher service account for function",
zap.String("service_account", fv1.FissionFetcherSA),
zap.String("service_account_namepsace", envNs),
zap.String("function_name", newFunc.ObjectMeta.Name),
zap.String("function_namespace", newFunc.ObjectMeta.Namespace))
}
}
},
})
return funcStore, controller
}
+2 -3
View File
@@ -71,8 +71,7 @@ type (
fissionClient *crd.FissionClient
fetcherConfig *fetcherConfig.Config
stopReadyPodControllerCh chan struct{}
readyPodController cache.Controller
readyPodIndexer cache.Indexer
readyPodInformer cache.SharedIndexInformer
readyPodQueue workqueue.DelayingInterface
poolInstanceID string // small random string to uniquify pod names
instanceID string // poolmgr instance id
@@ -249,7 +248,7 @@ func (gp *GenericPool) choosePod(newLabels map[string]string) (string, *apiv1.Po
key = item.(string)
gp.logger.Debug("got key from the queue", zap.String("key", key))
obj, exists, err := gp.readyPodIndexer.GetByKey(key)
obj, exists, err := gp.readyPodInformer.GetIndexer().GetByKey(key)
if err != nil {
gp.logger.Error("fetching object from store failed", zap.String("key", key), zap.Error(err))
return "", nil, err
+14 -14
View File
@@ -76,11 +76,10 @@ type (
enableIstio bool
fetcherConfig *fetcherConfig.Config
funcStore k8sCache.Store
funcController k8sCache.Controller
pkgStore k8sCache.Store
pkgController k8sCache.Controller
podInformer k8sCache.SharedIndexInformer
funcInformer *k8sCache.SharedIndexInformer
pkgInformer *k8sCache.SharedIndexInformer
podInformer k8sCache.SharedIndexInformer
defaultIdlePodReapTime time.Duration
}
@@ -103,7 +102,10 @@ func MakeGenericPoolManager(
metricsClient *metricsclient.Clientset,
functionNamespace string,
fetcherConfig *fetcherConfig.Config,
instanceID string) (executortype.ExecutorType, error) {
instanceID string,
funcInformer *k8sCache.SharedIndexInformer,
pkgInformer *k8sCache.SharedIndexInformer,
) (executortype.ExecutorType, error) {
gpmLogger := logger.Named("generic_pool_manager")
@@ -120,6 +122,8 @@ func MakeGenericPoolManager(
requestChannel: make(chan *request),
defaultIdlePodReapTime: 2 * time.Minute,
fetcherConfig: fetcherConfig,
funcInformer: funcInformer,
pkgInformer: pkgInformer,
}
go gpm.service()
@@ -132,16 +136,14 @@ func MakeGenericPoolManager(
gpm.enableIstio = istio
}
gpm.funcStore, gpm.funcController = gpm.makeFuncController(
gpm.fissionClient, gpm.kubernetesClient, gpm.namespace, gpm.enableIstio)
(*gpm.funcInformer).AddEventHandler(gpm.FunctionEventHandlers(gpm.kubernetesClient, gpm.namespace, gpm.enableIstio))
(*gpm.pkgInformer).AddEventHandler(gpm.PackageEventHandlers(gpm.kubernetesClient, gpm.namespace))
gpm.pkgStore, gpm.pkgController = gpm.makePkgController(gpm.fissionClient, gpm.kubernetesClient, gpm.namespace)
informerFactory, err := utils.GetInformerFacoryByExecutor(gpm.kubernetesClient, fv1.ExecutorTypePoolmgr)
kubeInformerFactory, err := utils.GetInformerFactoryByExecutor(gpm.kubernetesClient, fv1.ExecutorTypePoolmgr)
if err != nil {
return nil, err
}
gpm.podInformer = informerFactory.Core().V1().Pods().Informer()
gpm.podInformer = kubeInformerFactory.Core().V1().Pods().Informer()
return gpm, nil
}
@@ -149,8 +151,6 @@ func (gpm *GenericPoolManager) Run(ctx context.Context) {
// eagerPoolCreator must run after CleanupOldExecutorObjects.
// Otherwise, the poolmanager may wrongly delete the deployment.
go gpm.eagerPoolCreator()
go gpm.funcController.Run(ctx.Done())
go gpm.pkgController.Run(ctx.Done())
go gpm.podInformer.Run(ctx.Done())
go gpm.idleObjectReaper()
}
@@ -0,0 +1,99 @@
/*
Copyright 2018 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 (
"go.uber.org/zap"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/utils"
)
func (gpm *GenericPoolManager) PackageEventHandlers(kubernetesClient *kubernetes.Clientset, fissionfnNamespace string) k8sCache.ResourceEventHandlerFuncs {
return k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
pkg := obj.(*fv1.Package)
gpm.logger.Debug("list watch for package reported a new package addition",
zap.String("package_name", pkg.ObjectMeta.Name),
zap.String("package_namespace", pkg.ObjectMeta.Namespace))
// create or update role-binding for fetcher sa in env ns to be able to get the pkg contents from pkg namespace
envNs := fissionfnNamespace
if pkg.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = pkg.Spec.Environment.Namespace
}
// here, we return if we hit an error during rolebinding setup. this is because this rolebinding is mandatory for
// every function's package to be loaded into its env. without that, there's no point to move forward.
err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs)
if err != nil {
gpm.logger.Error("error creating rolebinding for package",
zap.Error(err),
zap.String("role_binding", fv1.PackageGetterRB),
zap.String("package_name", pkg.ObjectMeta.Name),
zap.String("package_namespace", pkg.ObjectMeta.Namespace))
return
}
gpm.logger.Debug("successfully set up rolebinding for fetcher service account",
zap.String("service_account", fv1.FissionFetcherSA),
zap.String("service_account_namespace", envNs),
zap.String("package_name", pkg.ObjectMeta.Name),
zap.String("package_namespace", pkg.ObjectMeta.Namespace))
},
UpdateFunc: func(oldObj, newObj interface{}) {
oldPkg := oldObj.(*fv1.Package)
newPkg := newObj.(*fv1.Package)
if oldPkg.ObjectMeta.ResourceVersion == newPkg.ObjectMeta.ResourceVersion {
return
}
// if a pkg's env reference gets updated and the newly referenced env is in a different ns,
// we need to update the role-binding in pkg ns to grant permissions to the fetcher-sa in env ns
// to do a get on pkg
if oldPkg.Spec.Environment.Namespace != newPkg.Spec.Environment.Namespace {
envNs := fissionfnNamespace
if newPkg.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = newPkg.Spec.Environment.Namespace
}
err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.PackageGetterRB,
newPkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole,
fv1.FissionFetcherSA, envNs)
if err != nil {
gpm.logger.Error("error updating rolebinding for package",
zap.Error(err),
zap.String("role_binding", fv1.PackageGetterRB),
zap.String("package_name", newPkg.ObjectMeta.Name),
zap.String("package_namespace", newPkg.ObjectMeta.Namespace))
return
}
gpm.logger.Debug("successfully updated rolebinding for fetcher service account",
zap.String("service_account", fv1.FissionFetcherSA),
zap.String("service_account_namespace", envNs),
zap.String("package_name", newPkg.ObjectMeta.Name),
zap.String("package_namespace", newPkg.ObjectMeta.Namespace))
}
},
}
}
@@ -1,111 +0,0 @@
/*
Copyright 2018 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 (
"time"
"go.uber.org/zap"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/fields"
"k8s.io/client-go/kubernetes"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/utils"
)
// TODO : It may make sense to make each of add, update, delete funcs run as separate go routines.
func (gpm *GenericPoolManager) makePkgController(fissionClient *crd.FissionClient,
kubernetesClient *kubernetes.Clientset, fissionfnNamespace string) (k8sCache.Store, k8sCache.Controller) {
resyncPeriod := 30 * time.Second
lw := k8sCache.NewListWatchFromClient(fissionClient.CoreV1().RESTClient(), "packages", metav1.NamespaceAll, fields.Everything())
pkgStore, controller := k8sCache.NewInformer(lw, &fv1.Package{}, resyncPeriod,
k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
pkg := obj.(*fv1.Package)
gpm.logger.Debug("list watch for package reported a new package addition",
zap.String("package_name", pkg.ObjectMeta.Name),
zap.String("package_namespace", pkg.ObjectMeta.Namespace))
// create or update role-binding for fetcher sa in env ns to be able to get the pkg contents from pkg namespace
envNs := fissionfnNamespace
if pkg.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = pkg.Spec.Environment.Namespace
}
// here, we return if we hit an error during rolebinding setup. this is because this rolebinding is mandatory for
// every function's package to be loaded into its env. without that, there's no point to move forward.
err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.PackageGetterRB, pkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole, fv1.FissionFetcherSA, envNs)
if err != nil {
gpm.logger.Error("error creating rolebinding for package",
zap.Error(err),
zap.String("role_binding", fv1.PackageGetterRB),
zap.String("package_name", pkg.ObjectMeta.Name),
zap.String("package_namespace", pkg.ObjectMeta.Namespace))
return
}
gpm.logger.Debug("successfully set up rolebinding for fetcher service account",
zap.String("service_account", fv1.FissionFetcherSA),
zap.String("service_account_namespace", envNs),
zap.String("package_name", pkg.ObjectMeta.Name),
zap.String("package_namespace", pkg.ObjectMeta.Namespace))
},
UpdateFunc: func(oldObj, newObj interface{}) {
oldPkg := oldObj.(*fv1.Package)
newPkg := newObj.(*fv1.Package)
if oldPkg.ObjectMeta.ResourceVersion == newPkg.ObjectMeta.ResourceVersion {
return
}
// if a pkg's env reference gets updated and the newly referenced env is in a different ns,
// we need to update the role-binding in pkg ns to grant permissions to the fetcher-sa in env ns
// to do a get on pkg
if oldPkg.Spec.Environment.Namespace != newPkg.Spec.Environment.Namespace {
envNs := fissionfnNamespace
if newPkg.Spec.Environment.Namespace != metav1.NamespaceDefault {
envNs = newPkg.Spec.Environment.Namespace
}
err := utils.SetupRoleBinding(gpm.logger, kubernetesClient, fv1.PackageGetterRB,
newPkg.ObjectMeta.Namespace, fv1.PackageGetterCR, fv1.ClusterRole,
fv1.FissionFetcherSA, envNs)
if err != nil {
gpm.logger.Error("error updating rolebinding for package",
zap.Error(err),
zap.String("role_binding", fv1.PackageGetterRB),
zap.String("package_name", newPkg.ObjectMeta.Name),
zap.String("package_namespace", newPkg.ObjectMeta.Namespace))
return
}
gpm.logger.Debug("successfully updated rolebinding for fetcher service account",
zap.String("service_account", fv1.FissionFetcherSA),
zap.String("service_account_namespace", envNs),
zap.String("package_name", newPkg.ObjectMeta.Name),
zap.String("package_namespace", newPkg.ObjectMeta.Namespace))
}
},
})
return pkgStore, controller
}
@@ -4,27 +4,31 @@ import (
"time"
"go.uber.org/zap"
apiv1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
informers "k8s.io/client-go/informers/core/v1"
"k8s.io/client-go/tools/cache"
"k8s.io/client-go/util/workqueue"
)
func (gp *GenericPool) newPodInformer() cache.SharedIndexInformer {
optionsModifier := func(options *metav1.ListOptions) {
options.LabelSelector = labels.Set(
gp.deployment.Spec.Selector.MatchLabels).AsSelector().String()
options.FieldSelector = "status.phase=Running"
}
return informers.NewFilteredPodInformer(gp.kubernetesClient, gp.namespace, 0, nil, optionsModifier)
}
func (gp *GenericPool) startReadyPodController() {
// create the pod watcher to filter by labels
// Filtering pod by phase=Running. In some cases the pod can be in
// different state than Running, for example Kubernetes sets a
// pod to Termination while k8s waits for the grace period of
// the pod, even if all the containers are in Ready state.
optionsModifier := func(options *metav1.ListOptions) {
options.LabelSelector = labels.Set(
gp.deployment.Spec.Selector.MatchLabels).AsSelector().String()
options.FieldSelector = "status.phase=Running"
}
readyPodWatcher := cache.NewFilteredListWatchFromClient(gp.kubernetesClient.CoreV1().RESTClient(), "pods", gp.namespace, optionsModifier)
gp.readyPodQueue = workqueue.NewDelayingQueue()
gp.readyPodIndexer, gp.readyPodController = cache.NewIndexerInformer(readyPodWatcher, &apiv1.Pod{}, 0, cache.ResourceEventHandlerFuncs{
gp.readyPodInformer = gp.newPodInformer()
gp.readyPodInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
key, err := cache.MetaNamespaceKeyFunc(obj)
if err == nil {
@@ -39,7 +43,7 @@ func (gp *GenericPool) startReadyPodController() {
gp.logger.Debug("delete func called", zap.String("key", key))
}
},
}, cache.Indexers{})
go gp.readyPodController.Run(gp.stopReadyPodControllerCh)
})
go gp.readyPodInformer.Run(gp.stopReadyPodControllerCh)
gp.logger.Info("readyPod controller started", zap.String("env", gp.env.ObjectMeta.Name), zap.String("envID", string(gp.env.ObjectMeta.UID)))
}
+8 -7
View File
@@ -33,8 +33,9 @@ type (
// reference into a resolveResult
functionReferenceResolver struct {
// FunctionReference -> function metadata
refCache *cache.Cache
store k8sCache.Store
refCache *cache.Cache
funcInformer *k8sCache.SharedIndexInformer
// store k8sCache.Store
}
resolveResultType int
@@ -68,10 +69,10 @@ const (
resolveResultMultipleFunctions
)
func makeFunctionReferenceResolver(store k8sCache.Store) *functionReferenceResolver {
func makeFunctionReferenceResolver(funcInformer *k8sCache.SharedIndexInformer) *functionReferenceResolver {
frr := &functionReferenceResolver{
refCache: cache.MakeCache(time.Minute, 0),
store: store,
refCache: cache.MakeCache(time.Minute, 0),
funcInformer: funcInformer,
}
return frr
}
@@ -120,7 +121,7 @@ func (frr *functionReferenceResolver) resolve(trigger fv1.HTTPTrigger) (*resolve
// resolveByName simply looks up function by name in a namespace.
func (frr *functionReferenceResolver) resolveByName(namespace, name string) (*resolveResult, error) {
// get function from cache
obj, isExist, err := frr.store.Get(&fv1.Function{
obj, isExist, err := (*frr.funcInformer).GetStore().Get(&fv1.Function{
ObjectMeta: metav1.ObjectMeta{
Namespace: namespace,
Name: name,
@@ -154,7 +155,7 @@ func (frr *functionReferenceResolver) resolveByFunctionWeights(namespace string,
for functionName, functionWeight := range fr.FunctionWeights {
// get function from cache
obj, isExist, err := frr.store.Get(&fv1.Function{
obj, isExist, err := (*frr.funcInformer).GetStore().Get(&fv1.Function{
ObjectMeta: metav1.ObjectMeta{
Namespace: namespace,
Name: functionName,
+70 -83
View File
@@ -23,8 +23,6 @@ import (
"github.com/gorilla/mux"
"go.uber.org/zap"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/fields"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
@@ -33,6 +31,7 @@ import (
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd"
executorClient "github.com/fission/fission/pkg/executor/client"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
"github.com/fission/fission/pkg/throttler"
"github.com/fission/fission/pkg/utils"
)
@@ -49,11 +48,9 @@ type HTTPTriggerSet struct {
resolver *functionReferenceResolver
crdClient rest.Interface
triggers []fv1.HTTPTrigger
triggerStore k8sCache.Store
triggerController k8sCache.Controller
triggerInformer k8sCache.SharedIndexInformer
functions []fv1.Function
funcStore k8sCache.Store
funcController k8sCache.Controller
funcInformer k8sCache.SharedIndexInformer
updateRouterRequestChannel chan struct{}
tsRoundTripperParams *tsRoundTripperParams
isDebugEnv bool
@@ -62,7 +59,7 @@ type HTTPTriggerSet struct {
}
func makeHTTPTriggerSet(logger *zap.Logger, fmap *functionServiceMap, fissionClient *crd.FissionClient,
kubeClient *kubernetes.Clientset, executor *executorClient.Client, crdClient rest.Interface, params *tsRoundTripperParams, isDebugEnv bool, unTapServiceTimeout time.Duration, actionThrottler *throttler.Throttler) (*HTTPTriggerSet, k8sCache.Store, k8sCache.Store) {
kubeClient *kubernetes.Clientset, executor *executorClient.Client, crdClient rest.Interface, params *tsRoundTripperParams, isDebugEnv bool, unTapServiceTimeout time.Duration, actionThrottler *throttler.Throttler) *HTTPTriggerSet {
httpTriggerSet := &HTTPTriggerSet{
logger: logger.Named("http_trigger_set"),
@@ -78,21 +75,19 @@ func makeHTTPTriggerSet(logger *zap.Logger, fmap *functionServiceMap, fissionCli
svcAddrUpdateThrottler: actionThrottler,
unTapServiceTimeout: unTapServiceTimeout,
}
var tStore, fnStore k8sCache.Store
var tController, fnController k8sCache.Controller
if httpTriggerSet.crdClient != nil {
tStore, tController = httpTriggerSet.initTriggerController()
httpTriggerSet.triggerStore = tStore
httpTriggerSet.triggerController = tController
fnStore, fnController = httpTriggerSet.initFunctionController()
httpTriggerSet.funcStore = fnStore
httpTriggerSet.funcController = fnController
informerFactory := genInformer.NewSharedInformerFactory(fissionClient, time.Second*30)
httpTriggerSet.triggerInformer = informerFactory.Core().V1().HTTPTriggers().Informer()
httpTriggerSet.funcInformer = informerFactory.Core().V1().Functions().Informer()
httpTriggerSet.addTriggerHandlers()
httpTriggerSet.addFunctionHandlers()
}
return httpTriggerSet, tStore, fnStore
return httpTriggerSet
}
func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter, resolver *functionReferenceResolver) {
func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter) {
resolver := makeFunctionReferenceResolver(&ts.funcInformer)
ts.resolver = resolver
ts.mutableRouter = mr
@@ -104,8 +99,8 @@ func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter
}
go ts.updateRouter()
go ts.syncTriggers()
go ts.runWatcher(ctx, ts.funcController)
go ts.runWatcher(ctx, ts.triggerController)
go ts.runInformer(ctx, ts.funcInformer)
go ts.runInformer(ctx, ts.triggerInformer)
}
func defaultHomeHandler(w http.ResponseWriter, r *http.Request) {
@@ -234,78 +229,70 @@ func (ts *HTTPTriggerSet) updateTriggerStatusFailed(ht *fv1.HTTPTrigger, err err
// TODO
}
func (ts *HTTPTriggerSet) initTriggerController() (k8sCache.Store, k8sCache.Controller) {
resyncPeriod := 30 * time.Second
listWatch := k8sCache.NewListWatchFromClient(ts.crdClient, "httptriggers", metav1.NamespaceAll, fields.Everything())
store, controller := k8sCache.NewInformer(listWatch, &fv1.HTTPTrigger{}, resyncPeriod,
k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
trigger := obj.(*fv1.HTTPTrigger)
go createIngress(ts.logger, trigger, ts.kubeClient)
ts.syncTriggers()
},
DeleteFunc: func(obj interface{}) {
ts.syncTriggers()
trigger := obj.(*fv1.HTTPTrigger)
go deleteIngress(ts.logger, trigger, ts.kubeClient)
},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
oldTrigger := oldObj.(*fv1.HTTPTrigger)
newTrigger := newObj.(*fv1.HTTPTrigger)
func (ts *HTTPTriggerSet) addTriggerHandlers() {
ts.triggerInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
trigger := obj.(*fv1.HTTPTrigger)
go createIngress(ts.logger, trigger, ts.kubeClient)
ts.syncTriggers()
},
DeleteFunc: func(obj interface{}) {
ts.syncTriggers()
trigger := obj.(*fv1.HTTPTrigger)
go deleteIngress(ts.logger, trigger, ts.kubeClient)
},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
oldTrigger := oldObj.(*fv1.HTTPTrigger)
newTrigger := newObj.(*fv1.HTTPTrigger)
if oldTrigger.ObjectMeta.ResourceVersion == newTrigger.ObjectMeta.ResourceVersion {
return
}
if oldTrigger.ObjectMeta.ResourceVersion == newTrigger.ObjectMeta.ResourceVersion {
return
}
go updateIngress(ts.logger, oldTrigger, newTrigger, ts.kubeClient)
ts.syncTriggers()
},
})
return store, controller
go updateIngress(ts.logger, oldTrigger, newTrigger, ts.kubeClient)
ts.syncTriggers()
},
})
}
func (ts *HTTPTriggerSet) initFunctionController() (k8sCache.Store, k8sCache.Controller) {
resyncPeriod := 30 * time.Second
listWatch := k8sCache.NewListWatchFromClient(ts.crdClient, "functions", metav1.NamespaceAll, fields.Everything())
store, controller := k8sCache.NewInformer(listWatch, &fv1.Function{}, resyncPeriod,
k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
ts.syncTriggers()
},
DeleteFunc: func(obj interface{}) {
ts.syncTriggers()
},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
oldFn := oldObj.(*fv1.Function)
fn := newObj.(*fv1.Function)
func (ts *HTTPTriggerSet) addFunctionHandlers() {
ts.funcInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
ts.syncTriggers()
},
DeleteFunc: func(obj interface{}) {
ts.syncTriggers()
},
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
oldFn := oldObj.(*fv1.Function)
fn := newObj.(*fv1.Function)
if oldFn.ObjectMeta.ResourceVersion == fn.ObjectMeta.ResourceVersion {
return
}
if oldFn.ObjectMeta.ResourceVersion == fn.ObjectMeta.ResourceVersion {
return
}
// update resolver function reference cache
for key, rr := range ts.resolver.copy() {
if key.namespace == fn.ObjectMeta.Namespace &&
rr.functionMap[fn.ObjectMeta.Name] != nil &&
rr.functionMap[fn.ObjectMeta.Name].ObjectMeta.ResourceVersion != fn.ObjectMeta.ResourceVersion {
// invalidate resolver cache
ts.logger.Debug("invalidating resolver cache")
err := ts.resolver.delete(key.namespace, key.triggerName, key.triggerResourceVersion)
if err != nil {
ts.logger.Error("error deleting functionReferenceResolver cache", zap.Error(err))
}
break
// update resolver function reference cache
for key, rr := range ts.resolver.copy() {
if key.namespace == fn.ObjectMeta.Namespace &&
rr.functionMap[fn.ObjectMeta.Name] != nil &&
rr.functionMap[fn.ObjectMeta.Name].ObjectMeta.ResourceVersion != fn.ObjectMeta.ResourceVersion {
// invalidate resolver cache
ts.logger.Debug("invalidating resolver cache")
err := ts.resolver.delete(key.namespace, key.triggerName, key.triggerResourceVersion)
if err != nil {
ts.logger.Error("error deleting functionReferenceResolver cache", zap.Error(err))
}
break
}
ts.syncTriggers()
},
})
return store, controller
}
ts.syncTriggers()
},
})
}
func (ts *HTTPTriggerSet) runWatcher(ctx context.Context, controller k8sCache.Controller) {
func (ts *HTTPTriggerSet) runInformer(ctx context.Context, informer k8sCache.SharedIndexInformer) {
go func() {
controller.Run(ctx.Done())
informer.Run(ctx.Done())
}()
}
@@ -316,7 +303,7 @@ func (ts *HTTPTriggerSet) syncTriggers() {
func (ts *HTTPTriggerSet) updateRouter() {
for range ts.updateRouterRequestChannel {
// get triggers
latestTriggers := ts.triggerStore.List()
latestTriggers := ts.triggerInformer.GetStore().List()
triggers := make([]fv1.HTTPTrigger, len(latestTriggers))
for _, t := range latestTriggers {
triggers = append(triggers, *t.(*fv1.HTTPTrigger))
@@ -324,7 +311,7 @@ func (ts *HTTPTriggerSet) updateRouter() {
ts.triggers = triggers
// get functions
latestFunctions := ts.funcStore.List()
latestFunctions := ts.funcInformer.GetStore().List()
functionTimeout := make(map[types.UID]int, len(latestFunctions))
functions := make([]fv1.Function, len(latestFunctions))
for _, f := range latestFunctions {
+6 -8
View File
@@ -63,7 +63,7 @@ import (
// request url ---[trigger]---> Function(name, deployment) ----[deployment]----> Function(name, uid) ----[pool mgr]---> k8s service url
func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTriggerSet, resolver *functionReferenceResolver) *mutableRouter {
func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTriggerSet) *mutableRouter {
var mr *mutableRouter
// see issue https://github.com/fission/fission/issues/1317
@@ -74,13 +74,13 @@ func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTrigger
mr = newMutableRouter(logger, mux.NewRouter())
}
httpTriggerSet.subscribeRouter(ctx, mr, resolver)
httpTriggerSet.subscribeRouter(ctx, mr)
return mr
}
func serve(ctx context.Context, logger *zap.Logger, port int, tracingSamplingRate float64,
httpTriggerSet *HTTPTriggerSet, resolver *functionReferenceResolver, displayAccessLog bool) {
mr := router(ctx, logger, httpTriggerSet, resolver)
httpTriggerSet *HTTPTriggerSet, displayAccessLog bool) {
mr := router(ctx, logger, httpTriggerSet)
url := fmt.Sprintf(":%v", port)
err := http.ListenAndServe(url, &ochttp.Handler{
@@ -236,7 +236,7 @@ func Start(logger *zap.Logger, port int, executorURL string) {
zap.Bool("default", displayAccessLog))
}
triggers, _, fnStore := makeHTTPTriggerSet(logger.Named("triggerset"), fmap, fissionClient, kubeClient, executor, fissionClient.CoreV1().RESTClient(), &tsRoundTripperParams{
triggers := makeHTTPTriggerSet(logger.Named("triggerset"), fmap, fissionClient, kubeClient, executor, fissionClient.CoreV1().RESTClient(), &tsRoundTripperParams{
timeout: timeout,
timeoutExponent: timeoutExponent,
disableKeepAlive: disableKeepAlive,
@@ -245,12 +245,10 @@ func Start(logger *zap.Logger, port int, executorURL string) {
svcAddrRetryCount: svcAddrRetryCount,
}, isDebugEnv, unTapServiceTimeout, throttler.MakeThrottler(svcAddrUpdateTimeout))
resolver := makeFunctionReferenceResolver(fnStore)
go serveMetric(logger)
logger.Info("starting router", zap.Int("port", port))
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
serve(ctx, logger, port, tracingSamplingRate, triggers, resolver, displayAccessLog)
serve(ctx, logger, port, tracingSamplingRate, triggers, displayAccessLog)
}
+16 -1
View File
@@ -1,17 +1,21 @@
package utils
import (
"fmt"
apiv1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
"k8s.io/apimachinery/pkg/selection"
k8sInformers "k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
k8sCache "k8s.io/client-go/tools/cache"
metricsapi "k8s.io/metrics/pkg/apis/metrics"
v1 "github.com/fission/fission/pkg/apis/core/v1"
)
func GetInformerFacoryByExecutor(client *kubernetes.Clientset, executorType v1.ExecutorType) (k8sInformers.SharedInformerFactory, error) {
func GetInformerFactoryByExecutor(client *kubernetes.Clientset, executorType v1.ExecutorType) (k8sInformers.SharedInformerFactory, error) {
executorLabel, err := labels.NewRequirement(v1.EXECUTOR_TYPE, selection.DoubleEquals, []string{string(executorType)})
if err != nil {
return nil, err
@@ -43,3 +47,14 @@ func SupportedMetricsAPIVersionAvailable(discoveredAPIGroups *metav1.APIGroupLis
}
return false
}
func GetCachedItem(obj apiv1.ObjectReference, informer k8sCache.SharedIndexInformer) (item interface{}, exists bool, err error) {
store := informer.GetStore()
item, exists, err = store.Get(obj)
if err != nil || !exists {
item, exists, err = store.GetByKey(fmt.Sprintf("%s/%s", obj.Namespace, obj.Name))
}
return item, exists, err
}
+1 -1
View File
@@ -84,7 +84,7 @@ main() {
$ROOT/test/tests/mqtrigger/nats/test_mqtrigger.sh \
$ROOT/test/tests/mqtrigger/nats/test_mqtrigger_error.sh \
$ROOT/test/tests/test_huge_response/test_huge_response.sh \
$ROOT/test/tests/test_kubectl/test_kubectl.sh
$ROOT/test/tests/test_kubectl/test_kubectl.sh \
$ROOT/test/tests/websocket/test_ws.sh
export JOBS=3