Add informers and internal go routines in manager (#2870)

* used manager's Add function in more places
* exit when ctx.Done is received in archivePruner go routines
* fix manager tests
* fix data race
* added more gpm function in manager and removed manager from a util function
* closed unused channel and stopped ticker after context is done
* added log statements
* used context.Done inside function instead of stopper channel
This commit is contained in:
Vardhaman Surana
2023-11-10 12:42:21 +05:30
committed by GitHub
parent 2a40b4538c
commit 3fabf64b3c
28 changed files with 195 additions and 114 deletions
+2 -2
View File
@@ -80,7 +80,7 @@ func Run(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *za
}
// do specialization in other goroutine to prevent blocking in newdeploy
go func() {
mgr.Add(ctx, func(_ context.Context) {
if *specializeOnStart {
var specializeReq fetcher.FunctionSpecializeRequest
@@ -95,7 +95,7 @@ func Run(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *za
}
}
atomic.StoreUint32(&readyToServe, 1)
}()
})
mux := http.NewServeMux()
mux.HandleFunc("/fetch", f.FetchHandler)
+12 -12
View File
@@ -52,8 +52,8 @@ func runWebhook(ctx context.Context, logger *zap.Logger, port int) error {
return webhook.Start(ctx, logger, port)
}
func runCanaryConfigServer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger) error {
return canaryconfigmgr.StartCanaryServer(ctx, clientGen, logger, false)
func runCanaryConfigServer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface) error {
return canaryconfigmgr.StartCanaryServer(ctx, clientGen, logger, mgr, false)
}
func runRouter(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, port int, executorUrl string) error {
@@ -64,12 +64,12 @@ func runExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface, lo
return executor.StartExecutor(ctx, clientGen, logger, mgr, port)
}
func runKubeWatcher(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerUrl string) error {
return kubewatcher.Start(ctx, clientGen, logger, routerUrl)
func runKubeWatcher(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerUrl string) error {
return kubewatcher.Start(ctx, clientGen, logger, mgr, routerUrl)
}
func runTimer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerUrl string) error {
return timer.Start(ctx, clientGen, logger, routerUrl)
func runTimer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerUrl string) error {
return timer.Start(ctx, clientGen, logger, mgr, routerUrl)
}
func runMessageQueueMgr(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerUrl string) error {
@@ -77,8 +77,8 @@ func runMessageQueueMgr(ctx context.Context, clientGen crd.ClientGeneratorInterf
}
// KEDA based MessageQueue Trigger Manager
func runMQManager(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerURL string) error {
return mqt.StartScalerManager(ctx, clientGen, logger, routerURL)
func runMQManager(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerURL string) error {
return mqt.StartScalerManager(ctx, clientGen, logger, mgr, routerURL)
}
func runStorageSvc(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, port int, storage storagesvc.Storage) error {
@@ -241,7 +241,7 @@ Options:
}
if arguments["--canaryConfig"] == true {
err := runCanaryConfigServer(ctx, clientGen, logger)
err := runCanaryConfigServer(ctx, clientGen, logger, mgr)
if err != nil {
logger.Error("canary config server exited with error: ", zap.Error(err))
return
@@ -267,7 +267,7 @@ Options:
}
if arguments["--kubewatcher"] == true {
err = runKubeWatcher(ctx, clientGen, logger, routerUrl)
err = runKubeWatcher(ctx, clientGen, logger, mgr, routerUrl)
if err != nil {
logger.Error("kubewatcher exited", zap.Error(err))
return
@@ -275,7 +275,7 @@ Options:
}
if arguments["--timer"] == true {
err = runTimer(ctx, clientGen, logger, routerUrl)
err = runTimer(ctx, clientGen, logger, mgr, routerUrl)
if err != nil {
logger.Error("timer exited", zap.Error(err))
return
@@ -291,7 +291,7 @@ Options:
}
if arguments["--mqt_keda"] == true {
err = runMQManager(ctx, clientGen, logger, routerUrl)
err = runMQManager(ctx, clientGen, logger, mgr, routerUrl)
if err != nil {
logger.Error("mqt scaler manager exited", zap.Error(err))
return
+1 -1
View File
@@ -61,10 +61,10 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *
}
envWatcher, err := makeEnvironmentWatcher(ctx, bmLogger, fissionClient, kubernetesClient, fetcherConfig, podSpecPatch)
envWatcher.Run(ctx)
if err != nil {
return err
}
envWatcher.Run(ctx, mgr)
pkgWatcher := makePackageWatcher(bmLogger, fissionClient,
kubernetesClient, storageSvcUrl,
+3 -4
View File
@@ -39,6 +39,7 @@ import (
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
"github.com/fission/fission/pkg/generated/clientset/versioned"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/manager"
)
const (
@@ -131,10 +132,8 @@ func (envw *environmentWatcher) getLabels(envName string, envNamespace string, e
}
}
func (envw *environmentWatcher) Run(ctx context.Context) {
for _, informer := range envw.envWatchInformer {
go informer.Run(ctx.Done())
}
func (envw *environmentWatcher) Run(ctx context.Context, mgr manager.Interface) {
mgr.AddInformers(ctx, envw.envWatchInformer)
}
func (envw *environmentWatcher) EnvWatchEventHandlers(ctx context.Context) error {
+2 -5
View File
@@ -314,18 +314,15 @@ func (pkgw *packageWatcher) Run(ctx context.Context, mgr manager.Interface) erro
mgr.Add(ctx, func(ctx context.Context) {
metrics.ServeMetrics(ctx, "buildermgr", pkgw.logger, mgr)
})
for _, podInformer := range pkgw.podInformer {
go podInformer.Run(ctx.Done())
}
mgr.AddInformers(ctx, pkgw.podInformer)
for _, pkgInformer := range pkgw.pkgInformer {
_, err := pkgInformer.AddEventHandler(pkgw.packageInformerHandler(ctx))
if err != nil {
pkgw.logger.Fatal("error adding package informer handler", zap.Error(err))
return err
}
go pkgInformer.Run(ctx.Done())
}
mgr.AddInformers(ctx, pkgw.pkgInformer)
return nil
}
+5 -6
View File
@@ -35,6 +35,7 @@ import (
"github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/generated/clientset/versioned"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/manager"
)
const (
@@ -136,10 +137,8 @@ func (canaryCfgMgr *canaryConfigMgr) CanaryConfigEventHandlers(ctx context.Conte
return nil
}
func (canaryCfgMgr *canaryConfigMgr) Run(ctx context.Context) {
for _, informer := range canaryCfgMgr.canaryConfigInformer {
go informer.Run(ctx.Done())
}
func (canaryCfgMgr *canaryConfigMgr) Run(ctx context.Context, mgr manager.Interface) {
mgr.AddInformers(ctx, canaryCfgMgr.canaryConfigInformer)
canaryCfgMgr.logger.Info("started canary configmgr controller")
}
@@ -558,7 +557,7 @@ func getEnvValue(envVar string) string {
return envVarSplit[1]
}
func StartCanaryServer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, unitTestFlag bool) error {
func StartCanaryServer(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, unitTestFlag bool) error {
cLogger := logger.Named("CanaryServer")
fissionClient, err := clientGen.GetFissionClient()
@@ -570,7 +569,7 @@ func StartCanaryServer(ctx context.Context, clientGen crd.ClientGeneratorInterfa
return fmt.Errorf("failed to get kubernetes client: %w", err)
}
err = ConfigureFeatures(ctx, cLogger, unitTestFlag, fissionClient, kubernetesClient)
err = ConfigureFeatures(ctx, cLogger, unitTestFlag, fissionClient, kubernetesClient, mgr)
if err != nil {
cLogger.Error("error configuring features - proceeding without optional features", zap.Error(err))
}
+4 -2
View File
@@ -25,10 +25,12 @@ import (
config "github.com/fission/fission/pkg/featureconfig"
"github.com/fission/fission/pkg/generated/clientset/versioned"
"github.com/fission/fission/pkg/utils/manager"
)
// ConfigureFeatures gets the feature config and configures the features that are enabled
func ConfigureFeatures(ctx context.Context, logger *zap.Logger, unitTestMode bool, fissionClient versioned.Interface, kubeClient kubernetes.Interface) error {
func ConfigureFeatures(ctx context.Context, logger *zap.Logger, unitTestMode bool, fissionClient versioned.Interface,
kubeClient kubernetes.Interface, mgr manager.Interface) error {
// set feature enabled to false if unitTestMode
if unitTestMode {
return nil
@@ -47,7 +49,7 @@ func ConfigureFeatures(ctx context.Context, logger *zap.Logger, unitTestMode boo
if err != nil {
return errors.Wrap(err, "failed to start canary config manager")
}
canaryCfgMgr.Run(ctx)
canaryCfgMgr.Run(ctx, mgr)
return err
}
+20 -11
View File
@@ -74,7 +74,7 @@ type (
)
// MakeExecutor returns an Executor for given ExecutorType(s).
func MakeExecutor(ctx context.Context, logger *zap.Logger, cms *cms.ConfigSecretController,
func MakeExecutor(ctx context.Context, logger *zap.Logger, mgr manager.Interface, cms *cms.ConfigSecretController,
fissionClient versioned.Interface, types map[fv1.ExecutorType]executortype.ExecutorType,
informers ...k8sCache.SharedIndexInformer) (*Executor, error) {
executor := &Executor{
@@ -88,17 +88,21 @@ func MakeExecutor(ctx context.Context, logger *zap.Logger, cms *cms.ConfigSecret
// Run all informers
for _, informer := range informers {
go informer.Run(ctx.Done())
informer := informer
mgr.Add(ctx, func(ctx context.Context) {
informer.Run(ctx.Done())
})
}
for _, et := range types {
go func(et executortype.ExecutorType) {
et.Run(ctx)
}(et)
et := et
mgr.Add(ctx, func(ctx context.Context) {
et.Run(ctx, mgr)
})
}
go executor.serveCreateFuncServices()
mgr.Add(ctx, func(ctx context.Context) {
executor.serveCreateFuncServices(ctx)
})
return executor, nil
}
@@ -108,9 +112,14 @@ func MakeExecutor(ctx context.Context, logger *zap.Logger, cms *cms.ConfigSecret
// get specialized. In other words, it ensures that when there's an
// ongoing request for a certain function, all other requests wait for
// that request to complete.
func (executor *Executor) serveCreateFuncServices() {
func (executor *Executor) serveCreateFuncServices(ctx context.Context) {
for {
req := <-executor.requestChan
var req *createFuncServiceRequest
select {
case <-ctx.Done():
return
case req = <-executor.requestChan:
}
function := req.function
fnName := k8sCache.MetaObjectToName(function)
fnkeyUR := crd.CacheKeyURFromObject(function)
@@ -384,7 +393,7 @@ func StartExecutor(ctx context.Context, clientGen crd.ClientGeneratorInterface,
informerFactory.Start(ctx.Done())
}
api, err := MakeExecutor(ctx, logger, cms, fissionClient, executorTypes,
api, err := MakeExecutor(ctx, logger, mgr, cms, fissionClient, executorTypes,
fissionInformers...,
)
if err != nil {
@@ -51,6 +51,7 @@ import (
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/manager"
"github.com/fission/fission/pkg/utils/maps"
otelUtils "github.com/fission/fission/pkg/utils/otel"
)
@@ -148,7 +149,7 @@ func MakeContainer(
}
// Run start the function along with an object reaper.
func (caaf *Container) Run(ctx context.Context) {
func (caaf *Container) Run(ctx context.Context, mgr manager.Interface) {
waitSynced := make([]k8sCache.InformerSynced, 0)
for _, deplListerSynced := range caaf.deplListerSynced {
waitSynced = append(waitSynced, deplListerSynced)
@@ -160,7 +161,9 @@ func (caaf *Container) Run(ctx context.Context) {
if ok := k8sCache.WaitForCacheSync(ctx.Done(), waitSynced...); !ok {
caaf.logger.Fatal("failed to wait for caches to sync")
}
go caaf.idleObjectReaper(ctx)
mgr.Add(ctx, func(ctx context.Context) {
caaf.idleObjectReaper(ctx)
})
}
// GetTypeName returns the executor type name.
+2 -1
View File
@@ -25,11 +25,12 @@ import (
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/executor/fscache"
"github.com/fission/fission/pkg/utils/manager"
)
type ExecutorType interface {
// Run runs background job.
Run(context.Context)
Run(context.Context, manager.Interface)
// GetTypeName returns the name of executor type
GetTypeName(context.Context) fv1.ExecutorType
@@ -53,6 +53,7 @@ import (
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/manager"
"github.com/fission/fission/pkg/utils/maps"
otelUtils "github.com/fission/fission/pkg/utils/otel"
)
@@ -162,7 +163,7 @@ func MakeNewDeploy(
}
// Run start the function and environment controller along with an object reaper.
func (deploy *NewDeploy) Run(ctx context.Context) {
func (deploy *NewDeploy) Run(ctx context.Context, mgr manager.Interface) {
waitSynced := make([]k8sCache.InformerSynced, 0)
for _, deplListerSynced := range deploy.deplListerSynced {
waitSynced = append(waitSynced, deplListerSynced)
@@ -174,7 +175,9 @@ func (deploy *NewDeploy) Run(ctx context.Context) {
if ok := k8sCache.WaitForCacheSync(ctx.Done(), waitSynced...); !ok {
deploy.logger.Fatal("failed to wait for caches to sync")
}
go deploy.idleObjectReaper(ctx)
mgr.Add(ctx, func(ctx context.Context) {
deploy.idleObjectReaper(ctx)
})
}
// GetTypeName returns the executor type name.
@@ -21,6 +21,7 @@ import (
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/loggerfactory"
"github.com/fission/fission/pkg/utils/manager"
"github.com/fission/fission/pkg/utils/uuid"
)
@@ -35,6 +36,8 @@ const (
func TestRefreshFuncPods(t *testing.T) {
os.Setenv("DEBUG_ENV", "true")
mgr := manager.New()
defer mgr.Wait()
logger := loggerfactory.GetLogger()
kubernetesClient := fake.NewSimpleClientset()
fissionClient := fClient.NewSimpleClientset()
@@ -70,7 +73,9 @@ func TestRefreshFuncPods(t *testing.T) {
}
ndm.nsResolver = &nsResolver
go ndm.Run(ctx)
mgr.Add(ctx, func(ctx context.Context) {
ndm.Run(ctx, mgr)
})
t.Log("New deploy manager started")
for _, f := range factory {
+23 -13
View File
@@ -53,6 +53,7 @@ import (
"github.com/fission/fission/pkg/generated/clientset/versioned"
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/manager"
otelUtils "github.com/fission/fission/pkg/utils/otel"
)
@@ -169,7 +170,7 @@ func MakeGenericPoolManager(ctx context.Context,
return gpm, nil
}
func (gpm *GenericPoolManager) Run(ctx context.Context) {
func (gpm *GenericPoolManager) Run(ctx context.Context, mgr manager.Interface) {
waitSynced := make([]k8sCache.InformerSynced, 0)
for _, podListerSynced := range gpm.podListerSynced {
waitSynced = append(waitSynced, podListerSynced)
@@ -179,10 +180,25 @@ func (gpm *GenericPoolManager) Run(ctx context.Context) {
}
go gpm.service()
gpm.poolPodC.InjectGpm(gpm)
go gpm.WebsocketStartEventChecker(ctx, gpm.kubernetesClient) //nolint:errcheck
go gpm.NoActiveConnectionEventChecker(ctx, gpm.kubernetesClient) //nolint:errcheck
go gpm.idleObjectReaper(ctx)
go gpm.poolPodC.Run(ctx, ctx.Done())
mgr.Add(ctx, func(ctx context.Context) {
err := gpm.WebsocketStartEventChecker(ctx, gpm.kubernetesClient)
if err != nil {
gpm.logger.Error("error in checking websocket start event from pod: ", zap.Error(err))
}
})
mgr.Add(ctx, func(ctx context.Context) {
err := gpm.NoActiveConnectionEventChecker(ctx, gpm.kubernetesClient) //nolint:errcheck
if err != nil {
gpm.logger.Error("error in checking inactive event from pod: ", zap.Error(err))
}
})
mgr.Add(ctx, func(ctx context.Context) {
gpm.idleObjectReaper(ctx)
})
mgr.Add(ctx, func(ctx context.Context) {
gpm.poolPodC.Run(ctx, ctx.Done(), mgr)
})
}
func (gpm *GenericPoolManager) GetTypeName(ctx context.Context) fv1.ExecutorType {
@@ -702,9 +718,6 @@ func (gpm *GenericPoolManager) doIdleObjectReaper(ctx context.Context) {
// WebsocketStartEventChecker checks if the pod has emitted a websocket connection start event
func (gpm *GenericPoolManager) WebsocketStartEventChecker(ctx context.Context, kubeClient kubernetes.Interface) error {
stopper := make(chan struct{})
defer close(stopper)
var wg wait.Group
for _, informer := range utils.GetInformerEventChecker(ctx, kubeClient, "WsConnectionStarted") {
_, err := informer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
@@ -727,7 +740,7 @@ func (gpm *GenericPoolManager) WebsocketStartEventChecker(ctx context.Context, k
if err != nil {
return err
}
wg.StartWithChannel(stopper, informer.Run)
wg.StartWithChannel(ctx.Done(), informer.Run)
}
wg.Wait()
return nil
@@ -735,9 +748,6 @@ func (gpm *GenericPoolManager) WebsocketStartEventChecker(ctx context.Context, k
// NoActiveConnectionEventChecker checks if the pod has emitted an inactive event
func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(ctx context.Context, kubeClient kubernetes.Interface) error {
stopper := make(chan struct{})
defer close(stopper)
var wg wait.Group
for _, informer := range utils.GetInformerEventChecker(ctx, kubeClient, "NoActiveConnections") {
_, err := informer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
@@ -771,7 +781,7 @@ func (gpm *GenericPoolManager) NoActiveConnectionEventChecker(ctx context.Contex
if err != nil {
return err
}
wg.StartWithChannel(stopper, informer.Run)
wg.StartWithChannel(ctx.Done(), informer.Run)
}
wg.Wait()
return nil
@@ -41,6 +41,7 @@ import (
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"
"github.com/fission/fission/pkg/utils/manager"
)
type (
@@ -246,7 +247,7 @@ func (p *PoolPodController) enqueueEnvDelete(obj interface{}) {
p.envDeleteQueue.Add(env)
}
func (p *PoolPodController) Run(ctx context.Context, stopCh <-chan struct{}) {
func (p *PoolPodController) Run(ctx context.Context, stopCh <-chan struct{}, mgr manager.Interface) {
defer utilruntime.HandleCrash()
defer p.envCreateUpdateQueue.ShutDown()
defer p.envDeleteQueue.ShutDown()
@@ -265,10 +266,16 @@ func (p *PoolPodController) Run(ctx context.Context, stopCh <-chan struct{}) {
p.logger.Fatal("failed to wait for caches to sync")
}
for i := 0; i < 4; i++ {
go wait.Until(p.workerRun(ctx, "envCreateUpdate", p.envCreateUpdateQueueProcessFunc), time.Second, stopCh)
mgr.Add(ctx, func(ctx context.Context) {
wait.Until(p.workerRun(ctx, "envCreateUpdate", p.envCreateUpdateQueueProcessFunc), time.Second, stopCh)
})
}
go wait.Until(p.workerRun(ctx, "envDeleteQueue", p.envDeleteQueueProcessFunc), time.Second, stopCh)
go wait.Until(p.workerRun(ctx, "spCleanupPodQueue", p.spCleanupPodQueueProcessFunc), time.Second, stopCh)
mgr.Add(ctx, func(ctx context.Context) {
wait.Until(p.workerRun(ctx, "envDeleteQueue", p.envDeleteQueueProcessFunc), time.Second, stopCh)
})
mgr.Add(ctx, func(ctx context.Context) {
wait.Until(p.workerRun(ctx, "spCleanupPodQueue", p.spCleanupPodQueueProcessFunc), time.Second, stopCh)
})
p.logger.Info("Started workers for poolPodController")
<-stopCh
p.logger.Info("Shutting down workers for poolPodController")
@@ -34,9 +34,12 @@ import (
genInformer "github.com/fission/fission/pkg/generated/informers/externalversions"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/loggerfactory"
"github.com/fission/fission/pkg/utils/manager"
)
func TestPoolPodControllerPodCleanup(t *testing.T) {
mgr := manager.New()
defer mgr.Wait()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
logger := loggerfactory.GetLogger()
@@ -74,7 +77,7 @@ func TestPoolPodControllerPodCleanup(t *testing.T) {
gpm := executor.(*GenericPoolManager)
ppc.InjectGpm(gpm)
go ppc.Run(ctx, ctx.Done())
go ppc.Run(ctx, ctx.Done(), mgr)
for _, f := range factory {
f.Start(ctx.Done())
+3 -2
View File
@@ -24,9 +24,10 @@ import (
"github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/publisher"
"github.com/fission/fission/pkg/utils/manager"
)
func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerUrl string) error {
func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerUrl string) error {
fissionClient, err := clientGen.GetFissionClient()
if err != nil {
return errors.Wrap(err, "failed to get fission client")
@@ -47,7 +48,7 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *
if err != nil {
return errors.Wrap(err, "error making watch sync")
}
ws.Run(ctx)
ws.Run(ctx, mgr)
return nil
}
+3 -4
View File
@@ -26,6 +26,7 @@ import (
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/generated/clientset/versioned"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/manager"
)
type (
@@ -51,10 +52,8 @@ func MakeWatchSync(ctx context.Context, logger *zap.Logger, client versioned.Int
return ws, nil
}
func (ws *WatchSync) Run(ctx context.Context) {
for _, informer := range ws.kubeWatcherInformer {
go informer.Run(ctx.Done())
}
func (ws *WatchSync) Run(ctx context.Context, mgr manager.Interface) {
mgr.AddInformers(ctx, ws.kubeWatcherInformer)
}
func (ws *WatchSync) KubeWatcherEventHandlers(ctx context.Context) error {
+3 -1
View File
@@ -86,7 +86,9 @@ func (mqt *MessageQueueTriggerManager) Run(ctx context.Context, mgr manager.Inte
if err != nil {
return err
}
go informer.Run(ctx.Done())
mgr.Add(ctx, func(ctx context.Context) {
informer.Run(ctx.Done())
})
if ok := k8sCache.WaitForCacheSync(ctx.Done(), informer.HasSynced); !ok {
mqt.logger.Fatal("failed to wait for caches to sync")
}
+5 -2
View File
@@ -25,6 +25,7 @@ import (
"github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/executor/util"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/manager"
)
var (
@@ -140,7 +141,7 @@ func mqTriggerEventHandlers(ctx context.Context, logger *zap.Logger, kubeClient
// StartScalerManager watches for changes in MessageQueueTrigger and,
// Based on changes, it Creates, Updates and Deletes Objects of Kind ScaledObjects, AuthenticationTriggers and Deployments
func StartScalerManager(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerURL string) error {
func StartScalerManager(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerURL string) error {
fissionClient, err := clientGen.GetFissionClient()
if err != nil {
return errors.Wrap(err, "failed to get fission client")
@@ -164,7 +165,9 @@ func StartScalerManager(ctx context.Context, clientGen crd.ClientGeneratorInterf
if err != nil {
return err
}
go informer.Run(ctx.Done())
mgr.Add(ctx, func(ctx context.Context) {
informer.Run(ctx.Done())
})
if ok := k8sCache.WaitForCacheSync(ctx.Done(), informer.HasSynced); !ok {
logger.Fatal("failed to wait for caches to sync")
}
+15 -13
View File
@@ -38,6 +38,7 @@ import (
"github.com/fission/fission/pkg/info"
"github.com/fission/fission/pkg/throttler"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/manager"
"github.com/fission/fission/pkg/utils/metrics"
)
@@ -93,7 +94,7 @@ func makeHTTPTriggerSet(logger *zap.Logger, fmap *functionServiceMap, fissionCli
return httpTriggerSet, nil
}
func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter) error {
func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mgr manager.Interface, mr *mutableRouter) error {
resolver := makeFunctionReferenceResolver(ts.logger, ts.funcInformer)
ts.resolver = resolver
ts.mutableRouter = mr
@@ -108,10 +109,12 @@ func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter
ts.logger.Info("skipping continuous trigger updates")
return nil
}
go ts.updateRouter()
go ts.syncTriggers()
go ts.runInformer(ctx, ts.funcInformer)
go ts.runInformer(ctx, ts.triggerInformer)
mgr.Add(ctx, func(ctx context.Context) {
ts.updateRouter(ctx)
})
ts.syncTriggers()
mgr.AddInformers(ctx, ts.funcInformer)
mgr.AddInformers(ctx, ts.triggerInformer)
return nil
}
@@ -378,20 +381,19 @@ func (ts *HTTPTriggerSet) addFunctionHandlers() error {
return nil
}
func (ts *HTTPTriggerSet) runInformer(ctx context.Context, informer map[string]k8sCache.SharedIndexInformer) {
for _, inf := range informer {
go inf.Run(ctx.Done())
}
}
func (ts *HTTPTriggerSet) syncTriggers() {
ts.syncDebouncer(func() {
ts.updateRouterRequestChannel <- struct{}{}
})
}
func (ts *HTTPTriggerSet) updateRouter() {
for range ts.updateRouterRequestChannel {
func (ts *HTTPTriggerSet) updateRouter(ctx context.Context) {
for {
select {
case <-ctx.Done():
return
case <-ts.updateRouterRequestChannel:
}
// get triggers
alltriggers := make([]fv1.HTTPTrigger, 0)
for _, triggerInformer := range ts.triggerInformer {
+3 -3
View File
@@ -64,7 +64,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) (*mutableRouter, error) {
func router(ctx context.Context, logger *zap.Logger, mgr manager.Interface, httpTriggerSet *HTTPTriggerSet) (*mutableRouter, error) {
var mr *mutableRouter
mux := mux.NewRouter()
mux.Use(metrics.HTTPMetricMiddleware)
@@ -80,7 +80,7 @@ func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTrigger
mr = newMutableRouter(logger, mux)
}
err = httpTriggerSet.subscribeRouter(ctx, mr)
err = httpTriggerSet.subscribeRouter(ctx, mgr, mr)
if err != nil {
return nil, err
}
@@ -89,7 +89,7 @@ func router(ctx context.Context, logger *zap.Logger, httpTriggerSet *HTTPTrigger
func serve(ctx context.Context, logger *zap.Logger, mgr manager.Interface, port int,
httpTriggerSet *HTTPTriggerSet, displayAccessLog bool) error {
mr, err := router(ctx, logger, httpTriggerSet)
mr, err := router(ctx, logger, mgr, httpTriggerSet)
if err != nil {
return errors.Wrap(err, "error making router")
}
+31 -15
View File
@@ -26,6 +26,7 @@ import (
"github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/generated/clientset/versioned"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/manager"
"github.com/pkg/errors"
)
@@ -55,17 +56,24 @@ func MakeArchivePruner(logger *zap.Logger, clientGen crd.ClientGeneratorInterfac
}
// pruneArchives listens to archiveChannel for archive ids that need to be deleted
func (pruner *ArchivePruner) pruneArchives() {
func (pruner *ArchivePruner) pruneArchives(ctx context.Context) {
pruner.logger.Debug("listening to archiveChannel to prune archives")
for archiveID := range pruner.archiveChan {
pruner.logger.Info("sending delete request for archive",
zap.String("archive_id", archiveID))
if err := pruner.stowClient.removeFileByID(archiveID); err != nil {
// logging the error and continuing with other deletions.
// hopefully this archive will be deleted in the next iteration.
pruner.logger.Error("ignoring error while deleting archive",
zap.Error(err),
for {
select {
case archiveID := <-pruner.archiveChan:
pruner.logger.Info("sending delete request for archive",
zap.String("archive_id", archiveID))
if err := pruner.stowClient.removeFileByID(archiveID); err != nil {
// logging the error and continuing with other deletions.
// hopefully this archive will be deleted in the next iteration.
pruner.logger.Error("ignoring error while deleting archive",
zap.Error(err),
zap.String("archive_id", archiveID))
}
case <-ctx.Done():
close(pruner.archiveChan)
pruner.logger.Info("stopped listening to archiveChannel to prune archives, context cancelled")
return
}
}
}
@@ -142,12 +150,20 @@ func (pruner *ArchivePruner) getOrphanArchives(ctx context.Context) {
// Start starts a go routine that listens to a channel for archive IDs that need to deleted.
// Also wakes up at regular intervals to make a list of archive IDs that need to be reaped
// and sends them over to the channel for deletion
func (pruner *ArchivePruner) Start(ctx context.Context) {
func (pruner *ArchivePruner) Start(ctx context.Context, mgr manager.Interface) {
ticker := time.NewTicker(pruner.pruneInterval * time.Minute)
go pruner.pruneArchives()
for range ticker.C {
// This method fetches unused archive IDs and sends them to archiveChannel for deletion
// silencing the errors, hoping they go away in next iteration.
pruner.getOrphanArchives(ctx)
mgr.Add(ctx, func(ctx context.Context) {
pruner.pruneArchives(ctx)
})
for {
select {
case <-ticker.C:
// This method fetches unused archive IDs and sends them to archiveChannel for deletion
// silencing the errors, hoping they go away in next iteration.
pruner.getOrphanArchives(ctx)
case <-ctx.Done():
ticker.Stop()
return
}
}
}
+3 -1
View File
@@ -316,7 +316,9 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *
if err != nil {
return errors.Wrap(err, "Error creating archivePruner")
}
go pruner.Start(ctx)
mgr.Add(ctx, func(ctx context.Context) {
pruner.Start(ctx, mgr)
})
}
logger.Info("storage service started")
+3 -2
View File
@@ -24,9 +24,10 @@ import (
"github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/publisher"
"github.com/fission/fission/pkg/utils/manager"
)
func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, routerUrl string) error {
func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *zap.Logger, mgr manager.Interface, routerUrl string) error {
fissionClient, err := clientGen.GetFissionClient()
if err != nil {
return errors.Wrap(err, "failed to get fission client")
@@ -42,6 +43,6 @@ func Start(ctx context.Context, clientGen crd.ClientGeneratorInterface, logger *
if err != nil {
return errors.Wrap(err, "error making timer sync")
}
timerSync.Run(ctx)
timerSync.Run(ctx, mgr)
return nil
}
+3 -4
View File
@@ -27,6 +27,7 @@ import (
"github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/generated/clientset/versioned"
"github.com/fission/fission/pkg/utils"
"github.com/fission/fission/pkg/utils/manager"
)
type (
@@ -52,10 +53,8 @@ func MakeTimerSync(ctx context.Context, logger *zap.Logger, fissionClient versio
return ws, nil
}
func (ws *TimerSync) Run(ctx context.Context) {
for _, informer := range ws.timeTriggerInformer {
go informer.Run(ctx.Done())
}
func (ws *TimerSync) Run(ctx context.Context, mgr manager.Interface) {
mgr.AddInformers(ctx, ws.timeTriggerInformer)
}
func (ws *TimerSync) AddUpdateTimeTrigger(timeTrigger *fv1.TimeTrigger) {
+13
View File
@@ -4,6 +4,8 @@ import (
"context"
"sync"
"time"
k8sCache "k8s.io/client-go/tools/cache"
)
var _ Interface = &GroupManager{}
@@ -15,6 +17,8 @@ type Interface interface {
// and will also remove the "function" from the list when it completes
Add(ctx context.Context, function func(context.Context))
AddInformers(ctx context.Context, informers map[string]k8sCache.SharedIndexInformer)
// Wait blocks the execution of the process until all the go routines in the manager are completed.
Wait()
@@ -39,6 +43,15 @@ func (g *GroupManager) Add(ctx context.Context, f func(context.Context)) {
}()
}
func (g *GroupManager) AddInformers(ctx context.Context, informers map[string]k8sCache.SharedIndexInformer) {
for _, informer := range informers {
informer := informer
g.Add(ctx, func(ctxArg context.Context) {
informer.Run(ctxArg.Done())
})
}
}
func (g *GroupManager) Wait() {
g.wg.Wait()
}
+2
View File
@@ -53,6 +53,8 @@ func TestAddAndWaitWithTimeout(t *testing.T) {
err := mgr.WaitWithTimeout(1 * time.Second)
require.NotNil(t, err, "manager WaitWithTimeout did not return an error when timeout exceeded")
mgr = New()
expectedValue := 11
mgr.Add(context.Background(), func(ctx context.Context) {
time.Sleep(100 * time.Millisecond)
+3
View File
@@ -10,9 +10,12 @@ import (
"k8s.io/client-go/kubernetes/fake"
"github.com/fission/fission/pkg/utils/loggerfactory"
"github.com/fission/fission/pkg/utils/manager"
)
func TestServiceAccountCheck(t *testing.T) {
mgr := manager.New()
defer mgr.Wait()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
kubernetesClient := fake.NewSimpleClientset()