K8s informer to work with specific namespaces for builder manager (#2649)
* watch informer for buildermgr in specific namepspaces * code review changes
This commit is contained in:
@@ -23,7 +23,6 @@ import (
|
|||||||
"github.com/pkg/errors"
|
"github.com/pkg/errors"
|
||||||
"go.uber.org/zap"
|
"go.uber.org/zap"
|
||||||
apiv1 "k8s.io/api/core/v1"
|
apiv1 "k8s.io/api/core/v1"
|
||||||
k8sInformers "k8s.io/client-go/informers"
|
|
||||||
|
|
||||||
fv1 "github.com/fission/fission/pkg/apis/core/v1"
|
fv1 "github.com/fission/fission/pkg/apis/core/v1"
|
||||||
"github.com/fission/fission/pkg/crd"
|
"github.com/fission/fission/pkg/crd"
|
||||||
@@ -65,10 +64,9 @@ func Start(ctx context.Context, logger *zap.Logger, storageSvcUrl string) error
|
|||||||
envWatcher := makeEnvironmentWatcher(ctx, bmLogger, fissionClient, kubernetesClient, fetcherConfig, podSpecPatch)
|
envWatcher := makeEnvironmentWatcher(ctx, bmLogger, fissionClient, kubernetesClient, fetcherConfig, podSpecPatch)
|
||||||
envWatcher.Run(ctx)
|
envWatcher.Run(ctx)
|
||||||
|
|
||||||
k8sInformerFactory := k8sInformers.NewSharedInformerFactory(kubernetesClient, time.Minute*30)
|
|
||||||
podInformer := k8sInformerFactory.Core().V1().Pods().Informer()
|
|
||||||
pkgWatcher := makePackageWatcher(bmLogger, fissionClient,
|
pkgWatcher := makePackageWatcher(bmLogger, fissionClient,
|
||||||
kubernetesClient, storageSvcUrl, podInformer,
|
kubernetesClient, storageSvcUrl,
|
||||||
|
utils.GetK8sInformersForNamespaces(kubernetesClient, time.Minute*30, fv1.Pods),
|
||||||
utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.PackagesResource))
|
utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.PackagesResource))
|
||||||
pkgWatcher.Run(ctx)
|
pkgWatcher.Run(ctx)
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -41,7 +41,7 @@ type (
|
|||||||
fissionClient versioned.Interface
|
fissionClient versioned.Interface
|
||||||
nsResolver *utils.NamespaceResolver
|
nsResolver *utils.NamespaceResolver
|
||||||
k8sClient kubernetes.Interface
|
k8sClient kubernetes.Interface
|
||||||
podInformer k8sCache.SharedIndexInformer
|
podInformer map[string]k8sCache.SharedIndexInformer
|
||||||
pkgInformer map[string]k8sCache.SharedIndexInformer
|
pkgInformer map[string]k8sCache.SharedIndexInformer
|
||||||
storageSvcUrl string
|
storageSvcUrl string
|
||||||
buildCache *cache.Cache
|
buildCache *cache.Cache
|
||||||
@@ -49,7 +49,7 @@ type (
|
|||||||
)
|
)
|
||||||
|
|
||||||
func makePackageWatcher(logger *zap.Logger, fissionClient versioned.Interface, k8sClientSet kubernetes.Interface,
|
func makePackageWatcher(logger *zap.Logger, fissionClient versioned.Interface, k8sClientSet kubernetes.Interface,
|
||||||
storageSvcUrl string, podInformer k8sCache.SharedIndexInformer,
|
storageSvcUrl string, podInformer,
|
||||||
pkgInformer map[string]k8sCache.SharedIndexInformer) *packageWatcher {
|
pkgInformer map[string]k8sCache.SharedIndexInformer) *packageWatcher {
|
||||||
pkgw := &packageWatcher{
|
pkgw := &packageWatcher{
|
||||||
logger: logger.Named("package_watcher"),
|
logger: logger.Named("package_watcher"),
|
||||||
@@ -124,6 +124,8 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) {
|
|||||||
|
|
||||||
// Create a new BackOff for health check on environment builder pod
|
// Create a new BackOff for health check on environment builder pod
|
||||||
healthCheckBackOff := utils.NewDefaultBackOff()
|
healthCheckBackOff := utils.NewDefaultBackOff()
|
||||||
|
builderNs := pkgw.nsResolver.GetBuilderNS(env.ObjectMeta.Namespace)
|
||||||
|
|
||||||
//if err != nil {
|
//if err != nil {
|
||||||
// pkgw.logger.Error("Unable to create BackOff for Health Check", zap.Error(err))
|
// pkgw.logger.Error("Unable to create BackOff for Health Check", zap.Error(err))
|
||||||
//}
|
//}
|
||||||
@@ -131,7 +133,7 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) {
|
|||||||
for healthCheckBackOff.NextExists() {
|
for healthCheckBackOff.NextExists() {
|
||||||
// Informer store is not able to use label to find the pod,
|
// Informer store is not able to use label to find the pod,
|
||||||
// iterate all available environment builders.
|
// iterate all available environment builders.
|
||||||
items := pkgw.podInformer.GetStore().List()
|
items := pkgw.podInformer[builderNs].GetStore().List()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
pkgw.logger.Error("error retrieving pod information for environment", zap.Error(err), zap.String("environment", env.ObjectMeta.Name))
|
pkgw.logger.Error("error retrieving pod information for environment", zap.Error(err), zap.String("environment", env.ObjectMeta.Name))
|
||||||
return
|
return
|
||||||
@@ -146,8 +148,6 @@ func (pkgw *packageWatcher) build(ctx context.Context, srcpkg *fv1.Package) {
|
|||||||
for _, item := range items {
|
for _, item := range items {
|
||||||
pod := item.(*apiv1.Pod)
|
pod := item.(*apiv1.Pod)
|
||||||
|
|
||||||
builderNs := pkgw.nsResolver.GetBuilderNS(env.ObjectMeta.Namespace)
|
|
||||||
|
|
||||||
// Filter non-matching pods
|
// Filter non-matching pods
|
||||||
if pod.ObjectMeta.Labels[LABEL_ENV_NAME] != env.ObjectMeta.Name ||
|
if pod.ObjectMeta.Labels[LABEL_ENV_NAME] != env.ObjectMeta.Name ||
|
||||||
pod.ObjectMeta.Labels[LABEL_ENV_NAMESPACE] != builderNs ||
|
pod.ObjectMeta.Labels[LABEL_ENV_NAMESPACE] != builderNs ||
|
||||||
@@ -327,7 +327,9 @@ func (pkgw *packageWatcher) packageInformerHandler(ctx context.Context) k8sCache
|
|||||||
|
|
||||||
func (pkgw *packageWatcher) Run(ctx context.Context) {
|
func (pkgw *packageWatcher) Run(ctx context.Context) {
|
||||||
go metrics.ServeMetrics(ctx, pkgw.logger)
|
go metrics.ServeMetrics(ctx, pkgw.logger)
|
||||||
go pkgw.podInformer.Run(ctx.Done())
|
for _, podInformer := range pkgw.podInformer {
|
||||||
|
go podInformer.Run(ctx.Done())
|
||||||
|
}
|
||||||
for _, pkgInformer := range pkgw.pkgInformer {
|
for _, pkgInformer := range pkgw.pkgInformer {
|
||||||
pkgInformer.AddEventHandler(pkgw.packageInformerHandler(ctx))
|
pkgInformer.AddEventHandler(pkgw.packageInformerHandler(ctx))
|
||||||
go pkgInformer.Run(ctx.Done())
|
go pkgInformer.Run(ctx.Done())
|
||||||
|
|||||||
Reference in New Issue
Block a user