112 lines
3.9 KiB
Go
112 lines
3.9 KiB
Go
// Package buildermgr — NSWatcher for multi-tenant mode.
|
|
//
|
|
// Listens for Namespaces labeled fission.io/managed=true and calls
|
|
// AddNamespace on envWatcher and packageWatcher so they pick up
|
|
// Environments and Packages in new tenant namespaces without a restart.
|
|
package buildermgr
|
|
|
|
import (
|
|
"context"
|
|
"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"
|
|
)
|
|
|
|
// StartNSWatcher watches for Namespaces labeled fission.io/managed=true
|
|
// and immediately registers per-NS informers in envWatcher and pkgWatcher.
|
|
func StartNSWatcher(
|
|
ctx context.Context,
|
|
logger *zap.Logger,
|
|
kubeClient kubernetes.Interface,
|
|
envw *environmentWatcher,
|
|
pkgw *packageWatcher,
|
|
mgr manager.Interface,
|
|
) {
|
|
nsManager := utils.NewNamespaceManager()
|
|
nsManager.Subscribe(NewNamespaceSubscriber(envw, pkgw, mgr))
|
|
if _, err := nsManager.BootstrapAndDispatch(ctx, utils.DefaultNSResolver().Snapshot(), utils.NamespaceSourceEnv, time.Now().UTC()); 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(k8sCache.ResourceEventHandlerFuncs{
|
|
AddFunc: func(obj interface{}) {
|
|
nsObj, ok := obj.(*corev1.Namespace)
|
|
if !ok || nsObj.Name == "" {
|
|
return
|
|
}
|
|
nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventAdd, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()))
|
|
if _, _, err := nsManager.DispatchAdd(ctx, nsObj.Name); err != nil {
|
|
logger.Error("buildermgr.NSWatcher: DispatchAdd failed",
|
|
zap.String("namespace", nsObj.Name), zap.Error(err))
|
|
}
|
|
},
|
|
UpdateFunc: func(oldObj, newObj interface{}) {
|
|
nsObj, ok := newObj.(*corev1.Namespace)
|
|
if !ok {
|
|
return
|
|
}
|
|
oldNSObj, _ := oldObj.(*corev1.Namespace)
|
|
if oldNSObj != nil && utils.IsManagedNamespace(oldNSObj.Labels) && !utils.IsManagedNamespace(nsObj.Labels) {
|
|
event := utils.NamespaceEventFromNamespace(utils.NamespaceEventRemove, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())
|
|
nsManager.Upsert(event)
|
|
logger.Info("buildermgr.NSWatcher: namespace removed from manager state; runtime registrations kept",
|
|
zap.String("namespace", nsObj.Name))
|
|
return
|
|
}
|
|
if !utils.IsManagedNamespace(nsObj.Labels) {
|
|
return
|
|
}
|
|
nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventUpdate, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()))
|
|
if _, _, err := nsManager.DispatchResync(ctx, nsObj.Name); err != nil {
|
|
logger.Error("buildermgr.NSWatcher: DispatchResync failed",
|
|
zap.String("namespace", nsObj.Name), zap.Error(err))
|
|
}
|
|
},
|
|
DeleteFunc: func(obj interface{}) {
|
|
event := utils.NamespaceEventFromObject(utils.NamespaceEventRemove, obj, utils.NamespaceSourceWatcher, time.Now().UTC())
|
|
if event.Name == "" {
|
|
return
|
|
}
|
|
nsManager.Upsert(event)
|
|
logger.Info("buildermgr.NSWatcher: namespace deleted from manager state; runtime registrations kept",
|
|
zap.String("namespace", event.Name))
|
|
},
|
|
})
|
|
|
|
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")
|
|
})
|
|
}
|
|
|
|
func builderNSName(obj interface{}) string {
|
|
nsObj, ok := obj.(*corev1.Namespace)
|
|
if !ok {
|
|
return ""
|
|
}
|
|
return nsObj.Name
|
|
}
|