From a1517ba4b21e6d501584b6ece81046bd761b557c Mon Sep 17 00:00:00 2001 From: Naeel Date: Sun, 26 Apr 2026 10:51:31 +0300 Subject: [PATCH] layer1: share managed namespace watcher startup --- .../2026-04-26-namespace-manager-step43.md | 19 ++++++++++++ pkg/buildermgr/ns_watcher.go | 27 +---------------- pkg/executor/multitenant/ns_watcher.go | 29 +------------------ pkg/router/ns_watcher.go | 27 +---------------- pkg/utils/namespace_manager.go | 27 +++++++++++++++++ 5 files changed, 49 insertions(+), 80 deletions(-) create mode 100644 doc/thinking/2026-04-26-namespace-manager-step43.md diff --git a/doc/thinking/2026-04-26-namespace-manager-step43.md b/doc/thinking/2026-04-26-namespace-manager-step43.md new file mode 100644 index 00000000..50e09c9b --- /dev/null +++ b/doc/thinking/2026-04-26-namespace-manager-step43.md @@ -0,0 +1,19 @@ +# 2026-04-26 — NamespaceManager rewrite, step 43 + +## Цель шага + +Убрать оставшуюся копипасту старта namespace informer-а из `buildermgr`, `router`, `executor/multitenant`. + +## Что меняем + +1. В `utils` добавляем `StartManagedNamespaceWatcher()`. +2. Helper централизует: + - informer factory с label selector; + - регистрацию event handlers; + - start/cache sync/stop logging через `mgr`. +3. Три watcher-а переходят на общий helper. + +## Что НЕ меняем + +- не меняем lifecycle logic; +- не меняем selector contract `fission.io/managed=true`. \ No newline at end of file diff --git a/pkg/buildermgr/ns_watcher.go b/pkg/buildermgr/ns_watcher.go index 84bb0807..844ae9bc 100644 --- a/pkg/buildermgr/ns_watcher.go +++ b/pkg/buildermgr/ns_watcher.go @@ -10,11 +10,7 @@ import ( "time" "go.uber.org/zap" - corev1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - k8sInformers "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" - k8sCache "k8s.io/client-go/tools/cache" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/manager" @@ -34,26 +30,5 @@ func StartNSWatcher( if err != nil { logger.Error("buildermgr.NSWatcher: BootstrapAndDispatch failed", zap.Error(err)) } - - factory := k8sInformers.NewSharedInformerFactoryWithOptions( - kubeClient, - 30*time.Minute, - k8sInformers.WithTweakListOptions(func(opts *metav1.ListOptions) { - opts.LabelSelector = utils.ManagedNamespaceLabelSelector() - }), - ) - - nsInformer := factory.Core().V1().Namespaces().Informer() - - _, _ = nsInformer.AddEventHandler(utils.NewNamespaceWatcherEventHandlers(ctx, logger, "buildermgr.NSWatcher", nsManager, utils.NamespaceRemovalStrategyTrackOnly)) - - mgr.Add(ctx, func(ctx context.Context) { - logger.Info("buildermgr.NSWatcher: started", - zap.String("label", utils.ManagedNamespaceLabelSelector())) - factory.Start(ctx.Done()) - factory.WaitForCacheSync(ctx.Done()) - logger.Info("buildermgr.NSWatcher: cache synced — watching for new namespaces") - <-ctx.Done() - logger.Info("buildermgr.NSWatcher: stopped") - }) + utils.StartManagedNamespaceWatcher(ctx, logger, "buildermgr.NSWatcher", kubeClient, mgr, utils.NewNamespaceWatcherEventHandlers(ctx, logger, "buildermgr.NSWatcher", nsManager, utils.NamespaceRemovalStrategyTrackOnly)) } diff --git a/pkg/executor/multitenant/ns_watcher.go b/pkg/executor/multitenant/ns_watcher.go index d55fadca..79c1539a 100644 --- a/pkg/executor/multitenant/ns_watcher.go +++ b/pkg/executor/multitenant/ns_watcher.go @@ -65,11 +65,7 @@ import ( "time" "go.uber.org/zap" - corev1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - k8sInformers "k8s.io/client-go/informers" "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/executor/executortype" @@ -96,30 +92,7 @@ func StartNSWatcher( if err != nil { logger.Error("multitenant.NSWatcher: BootstrapAndDispatch failed", zap.Error(err)) } - - // Use a label-filtered informer so only Namespaces with our label are delivered. - // The resync period of 30m is standard for Fission informers — it re-lists to recover - // from any missed events, but normal operation is purely event-driven (no ticking). - factory := k8sInformers.NewSharedInformerFactoryWithOptions( - kubernetesClient, - 30*time.Minute, - k8sInformers.WithTweakListOptions(func(opts *metav1.ListOptions) { - opts.LabelSelector = utils.ManagedNamespaceLabelSelector() - }), - ) - - nsInformer := factory.Core().V1().Namespaces().Informer() - - _, _ = nsInformer.AddEventHandler(utils.NewNamespaceWatcherEventHandlers(ctx, logger, "multitenant.NSWatcher", nsManager, utils.NamespaceRemovalStrategyTrackOnly)) - - mgr.Add(ctx, func(ctx context.Context) { - logger.Info("multitenant.NSWatcher: started", zap.String("label", utils.ManagedNamespaceLabelSelector())) - factory.Start(ctx.Done()) - factory.WaitForCacheSync(ctx.Done()) - logger.Info("multitenant.NSWatcher: cache synced — watching for new namespaces") - <-ctx.Done() - logger.Info("multitenant.NSWatcher: stopped") - }) + utils.StartManagedNamespaceWatcher(ctx, logger, "multitenant.NSWatcher", kubernetesClient, mgr, utils.NewNamespaceWatcherEventHandlers(ctx, logger, "multitenant.NSWatcher", nsManager, utils.NamespaceRemovalStrategyTrackOnly)) } // registerNamespace calls AddNamespace on every executor type for the given namespace. diff --git a/pkg/router/ns_watcher.go b/pkg/router/ns_watcher.go index 9bbd3f13..35f903fc 100644 --- a/pkg/router/ns_watcher.go +++ b/pkg/router/ns_watcher.go @@ -10,11 +10,7 @@ import ( "time" "go.uber.org/zap" - corev1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - k8sInformers "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" - k8sCache "k8s.io/client-go/tools/cache" "github.com/fission/fission/pkg/utils" "github.com/fission/fission/pkg/utils/manager" @@ -34,26 +30,5 @@ func StartNSWatcher( if err != nil { logger.Error("router.NSWatcher: BootstrapAndDispatch failed", zap.Error(err)) } - - factory := k8sInformers.NewSharedInformerFactoryWithOptions( - kubeClient, - 30*time.Minute, - k8sInformers.WithTweakListOptions(func(opts *metav1.ListOptions) { - opts.LabelSelector = utils.ManagedNamespaceLabelSelector() - }), - ) - - nsInformer := factory.Core().V1().Namespaces().Informer() - - _, _ = nsInformer.AddEventHandler(utils.NewNamespaceWatcherEventHandlers(ctx, logger, "router.NSWatcher", nsManager, utils.NamespaceRemovalStrategyTrackOnly)) - - mgr.Add(ctx, func(ctx context.Context) { - logger.Info("router.NSWatcher: started", - zap.String("label", utils.ManagedNamespaceLabelSelector())) - factory.Start(ctx.Done()) - factory.WaitForCacheSync(ctx.Done()) - logger.Info("router.NSWatcher: cache synced — watching for new namespaces") - <-ctx.Done() - logger.Info("router.NSWatcher: stopped") - }) + utils.StartManagedNamespaceWatcher(ctx, logger, "router.NSWatcher", kubeClient, mgr, utils.NewNamespaceWatcherEventHandlers(ctx, logger, "router.NSWatcher", nsManager, utils.NamespaceRemovalStrategyTrackOnly)) } diff --git a/pkg/utils/namespace_manager.go b/pkg/utils/namespace_manager.go index a97f00bf..36a49c90 100644 --- a/pkg/utils/namespace_manager.go +++ b/pkg/utils/namespace_manager.go @@ -9,7 +9,12 @@ import ( "go.uber.org/zap" corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + k8sInformers "k8s.io/client-go/informers" + "k8s.io/client-go/kubernetes" k8sCache "k8s.io/client-go/tools/cache" + + managerPkg "github.com/fission/fission/pkg/utils/manager" ) type NamespaceSubscriber interface { @@ -207,6 +212,28 @@ func NewNamespaceWatcherEventHandlers(ctx context.Context, logger *zap.Logger, c } } +func StartManagedNamespaceWatcher(ctx context.Context, logger *zap.Logger, component string, kubeClient kubernetes.Interface, mgr managerPkg.Interface, handlers k8sCache.ResourceEventHandlerFuncs) { + factory := k8sInformers.NewSharedInformerFactoryWithOptions( + kubeClient, + 30*time.Minute, + k8sInformers.WithTweakListOptions(func(opts *metav1.ListOptions) { + opts.LabelSelector = ManagedNamespaceLabelSelector() + }), + ) + + nsInformer := factory.Core().V1().Namespaces().Informer() + _, _ = nsInformer.AddEventHandler(handlers) + + mgr.Add(ctx, func(ctx context.Context) { + logger.Info(component+": started", zap.String("label", ManagedNamespaceLabelSelector())) + factory.Start(ctx.Done()) + factory.WaitForCacheSync(ctx.Done()) + logger.Info(component + ": cache synced — watching for new namespaces") + <-ctx.Done() + logger.Info(component + ": stopped") + }) +} + func (m *inMemoryNamespaceManager) Subscribe(subscriber NamespaceSubscriber) { m.mu.Lock() defer m.mu.Unlock()