layer1: share managed namespace watcher startup
This commit is contained in:
@@ -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`.
|
||||||
@@ -10,11 +10,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"go.uber.org/zap"
|
"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"
|
"k8s.io/client-go/kubernetes"
|
||||||
k8sCache "k8s.io/client-go/tools/cache"
|
|
||||||
|
|
||||||
"github.com/fission/fission/pkg/utils"
|
"github.com/fission/fission/pkg/utils"
|
||||||
"github.com/fission/fission/pkg/utils/manager"
|
"github.com/fission/fission/pkg/utils/manager"
|
||||||
@@ -34,26 +30,5 @@ func StartNSWatcher(
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Error("buildermgr.NSWatcher: BootstrapAndDispatch failed", zap.Error(err))
|
logger.Error("buildermgr.NSWatcher: BootstrapAndDispatch failed", zap.Error(err))
|
||||||
}
|
}
|
||||||
|
utils.StartManagedNamespaceWatcher(ctx, logger, "buildermgr.NSWatcher", kubeClient, mgr, utils.NewNamespaceWatcherEventHandlers(ctx, logger, "buildermgr.NSWatcher", nsManager, utils.NamespaceRemovalStrategyTrackOnly))
|
||||||
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")
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -65,11 +65,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"go.uber.org/zap"
|
"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"
|
"k8s.io/client-go/kubernetes"
|
||||||
k8sCache "k8s.io/client-go/tools/cache"
|
|
||||||
|
|
||||||
fv1 "github.com/fission/fission/pkg/apis/core/v1"
|
fv1 "github.com/fission/fission/pkg/apis/core/v1"
|
||||||
"github.com/fission/fission/pkg/executor/executortype"
|
"github.com/fission/fission/pkg/executor/executortype"
|
||||||
@@ -96,30 +92,7 @@ func StartNSWatcher(
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Error("multitenant.NSWatcher: BootstrapAndDispatch failed", zap.Error(err))
|
logger.Error("multitenant.NSWatcher: BootstrapAndDispatch failed", zap.Error(err))
|
||||||
}
|
}
|
||||||
|
utils.StartManagedNamespaceWatcher(ctx, logger, "multitenant.NSWatcher", kubernetesClient, mgr, utils.NewNamespaceWatcherEventHandlers(ctx, logger, "multitenant.NSWatcher", nsManager, utils.NamespaceRemovalStrategyTrackOnly))
|
||||||
// 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")
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// registerNamespace calls AddNamespace on every executor type for the given namespace.
|
// registerNamespace calls AddNamespace on every executor type for the given namespace.
|
||||||
|
|||||||
@@ -10,11 +10,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"go.uber.org/zap"
|
"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"
|
"k8s.io/client-go/kubernetes"
|
||||||
k8sCache "k8s.io/client-go/tools/cache"
|
|
||||||
|
|
||||||
"github.com/fission/fission/pkg/utils"
|
"github.com/fission/fission/pkg/utils"
|
||||||
"github.com/fission/fission/pkg/utils/manager"
|
"github.com/fission/fission/pkg/utils/manager"
|
||||||
@@ -34,26 +30,5 @@ func StartNSWatcher(
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
logger.Error("router.NSWatcher: BootstrapAndDispatch failed", zap.Error(err))
|
logger.Error("router.NSWatcher: BootstrapAndDispatch failed", zap.Error(err))
|
||||||
}
|
}
|
||||||
|
utils.StartManagedNamespaceWatcher(ctx, logger, "router.NSWatcher", kubeClient, mgr, utils.NewNamespaceWatcherEventHandlers(ctx, logger, "router.NSWatcher", nsManager, utils.NamespaceRemovalStrategyTrackOnly))
|
||||||
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")
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -9,7 +9,12 @@ import (
|
|||||||
|
|
||||||
"go.uber.org/zap"
|
"go.uber.org/zap"
|
||||||
corev1 "k8s.io/api/core/v1"
|
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"
|
k8sCache "k8s.io/client-go/tools/cache"
|
||||||
|
|
||||||
|
managerPkg "github.com/fission/fission/pkg/utils/manager"
|
||||||
)
|
)
|
||||||
|
|
||||||
type NamespaceSubscriber interface {
|
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) {
|
func (m *inMemoryNamespaceManager) Subscribe(subscriber NamespaceSubscriber) {
|
||||||
m.mu.Lock()
|
m.mu.Lock()
|
||||||
defer m.mu.Unlock()
|
defer m.mu.Unlock()
|
||||||
|
|||||||
Reference in New Issue
Block a user