From 4eedf95f5c120d4c5f1f9df8984fa79da14da443 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E2=80=9CNaeel=E2=80=9D?= Date: Mon, 18 May 2026 09:04:13 +0400 Subject: [PATCH] fix(namespace): executor/router/buildermgr RemoveNamespace + per-NS informer lifecycle - Add RemoveNamespace(ctx, ns) to executortype.ExecutorType interface - Implement RemoveNamespace in poolmgr, newdeploy, container executor types - Add per-namespace context cancellation (nsCancels map) in all three types so informer factories are stopped when namespace is removed (fixes goroutine leak) - Add PoolPodController.RemoveNamespace to clear envLister/podLister maps - Add deregisterNamespace() in executor multitenant subscriber - Switch executor/router/buildermgr watcher strategy from TrackOnly to DispatchRemove so RemoveFunc is called when fission.io/managed label is removed - Add RemoveFunc to executor/router/buildermgr namespace subscribers - Add RemoveNamespace to environmentWatcher and packageWatcher with per-NS cancel - Add RemoveNamespace to HTTPTriggerSet: cancels informers, removes from maps, calls syncTriggers - Fix ns_watcher_test.go fakeExecutorType to implement new RemoveNamespace method Fixes: - Executor dedup gap: re-added namespace was silently skipped (envLister/deplLister still present) - Goroutine/FD leak: old informer factories ran forever after namespace removal - Router stale routes: HTTPTriggers for removed namespace stayed in routing table --- .../FORENSIC_ARCHITECTURE_AUDIT.md | 0 doc/console-compat-2026-05-15.md | 72 +++++++++++++++++++ pkg/buildermgr/envwatcher.go | 21 +++++- pkg/buildermgr/namespace_subscriber.go | 25 +++++++ pkg/buildermgr/ns_watcher.go | 4 +- pkg/buildermgr/pkgwatcher.go | 23 +++++- .../executortype/container/containermgr.go | 38 ++++++++-- pkg/executor/executortype/executortype.go | 5 ++ .../executortype/newdeploy/newdeploymgr.go | 38 ++++++++-- pkg/executor/executortype/poolmgr/gpm.go | 39 +++++++++- .../executortype/poolmgr/poolpodcontroller.go | 11 +++ .../multitenant/namespace_subscriber.go | 3 + pkg/executor/multitenant/ns_watcher.go | 27 ++++++- pkg/executor/multitenant/ns_watcher_test.go | 1 + pkg/router/httpTriggers.go | 55 ++++++++++---- pkg/router/namespace_subscriber.go | 10 +++ pkg/router/ns_watcher.go | 4 +- 17 files changed, 347 insertions(+), 29 deletions(-) rename FORENSIC_ARCHITECTURE_AUDIT.md => doc/FORENSIC_ARCHITECTURE_AUDIT.md (100%) create mode 100644 doc/console-compat-2026-05-15.md diff --git a/FORENSIC_ARCHITECTURE_AUDIT.md b/doc/FORENSIC_ARCHITECTURE_AUDIT.md similarity index 100% rename from FORENSIC_ARCHITECTURE_AUDIT.md rename to doc/FORENSIC_ARCHITECTURE_AUDIT.md diff --git a/doc/console-compat-2026-05-15.md b/doc/console-compat-2026-05-15.md new file mode 100644 index 00000000..be3d442a --- /dev/null +++ b/doc/console-compat-2026-05-15.md @@ -0,0 +1,72 @@ +# Console ↔ fission-src multitenant: compatibility check (2026-05-15) + +## Что изменилось в fission-src (feature/multitenant) + +| Изменение | Файл | +|---|---| +| Удалён старый partial RBAC | `deploy/executor-ns-watcher-rbac.yaml` | +| Добавлен полный RBAC для NSWatcher | `deploy/multitenant/rbac.yaml` | +| Добавлен `EnsureNamespaceSA` | `pkg/utils/serviceaccount.go` | +| NSWatcher вызывает `EnsureNamespaceSA` при обнаружении NS с `fission.io/managed=true` | `pkg/executor/multitenant/ns_watcher.go` | +| Добавлен тест NSWatcher с fake k8s | `pkg/utils/namespace_manager_test.go` | + +--- + +## Что делает консоль при создании namespace + +`SetupFissionNamespace` в `console/internal/fission/namespace.go`: + +1. Создаёт Namespace с лейблами: + - `managed-by=fission-console` + - `fission.io/managed=true` ← триггер для NSWatcher + +2. Создаёт ServiceAccounts: `fission-fetcher`, `fission-builder` + +3. Создаёт RoleBindings с `cluster-admin` ClusterRole для всех Fission SA: + - `fission-executor`, `fission-router`, `fission-buildermgr`, `fission-kubewatcher`, `fission-timer` + - `fission-fetcher` (из fission NS + локально в user NS) + - `fission-builder` (из fission NS + локально в user NS) + +--- + +## Взаимодействие с EnsureNamespaceSA + +`EnsureNamespaceSA` вызывается NSWatcher **после** того как консоль создала NS. +Логика (в `setupSAAndRoleBindings`): + +1. Создаёт/получает SA `fission-fetcher` → SA уже существует → `IsAlreadyExists` → OK +2. Для каждого permission из `fetcherCheck` вызывает `checkPermission` через `localsubjectaccessreviews` +3. Поскольку у `fission-fetcher` уже есть `cluster-admin` RoleBinding (создан консолью) → + **все проверки возвращают `exists=true`** → `rules` остаётся пустым → Role и RoleBinding **не создаются** + +**Итог: EnsureNamespaceSA является no-op если консоль уже настроила namespace. Никаких конфликтов.** + +--- + +## Что нужно на кластере + +Для работы NSWatcher нужен `deploy/multitenant/rbac.yaml` применён **один раз**: +```bash +kubectl apply -f ~/terra/fission-src/deploy/multitenant/rbac.yaml +``` + +Это даёт: +- `fission-executor` → `list/watch namespaces` (NSWatcher) +- `fission-router` → `list/watch namespaces` (NSWatcher) +- `fission-executor` → `create SA/Role/RoleBinding` в user NS (`fission-executor-sa-provisioner`) + +Без этого RBAC `EnsureNamespaceSA` будет падать с Forbidden, но **консоль продолжит работать** — она создаёт SA/RoleBindings сама и не зависит от NSWatcher. + +--- + +## Вердикт + +| Сценарий | Статус | +|---|---| +| Новый NS создаётся через консоль | ✅ работает как раньше | +| NSWatcher обнаруживает NS по `fission.io/managed=true` | ✅ совместимо | +| `EnsureNamespaceSA` вызывается в уже настроенном NS | ✅ no-op, нет конфликтов | +| Старый NS (без нового fission-bundle) | ✅ консоль не зависит от NSWatcher | +| Сборка консоли (`go build ./...`) | ✅ BUILD OK | + +**Код консоли менять не нужно.** Нужно только применить `deploy/multitenant/rbac.yaml` при деплое нового fission-bundle. diff --git a/pkg/buildermgr/envwatcher.go b/pkg/buildermgr/envwatcher.go index 8da7075c..9567b0ab 100644 --- a/pkg/buildermgr/envwatcher.go +++ b/pkg/buildermgr/envwatcher.go @@ -76,6 +76,8 @@ type ( podSpecPatch *apiv1.PodSpec envWatchInformer map[string]k8sCache.SharedIndexInformer enableOwnerReferences bool + // nsCancels holds per-namespace context cancel functions. + nsCancels map[string]context.CancelFunc } ) @@ -111,8 +113,8 @@ func makeEnvironmentWatcher( podSpecPatch: podSpecPatch, envWatchInformer: utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.EnvironmentResource), enableOwnerReferences: utils.IsOwnerReferencesEnabled(), + nsCancels: make(map[string]context.CancelFunc), } - err := envWatcher.EnvWatchEventHandlers(ctx) if err != nil { return nil, err @@ -540,6 +542,21 @@ func (envw *environmentWatcher) AddNamespace(ctx context.Context, ns string, mgr envw.envWatchInformer[ns] = envInf mgr.AddInformers(ctx, map[string]k8sCache.SharedIndexInformer{ns: envInf}) - factory.Start(ctx.Done()) + + // Create a per-namespace cancellable context for informer lifecycle. + nsCtx, nsCancel := context.WithCancel(ctx) + envw.nsCancels[ns] = nsCancel + + factory.Start(nsCtx.Done()) envw.logger.Info("buildermgr.envWatcher.AddNamespace: done", zap.String("namespace", ns)) } + +// RemoveNamespace deregisters a namespace from the environment watcher. +func (envw *environmentWatcher) RemoveNamespace(ns string) { + if cancel, ok := envw.nsCancels[ns]; ok { + cancel() + delete(envw.nsCancels, ns) + } + delete(envw.envWatchInformer, ns) + envw.logger.Info("buildermgr.envWatcher.RemoveNamespace: cleaned up", zap.String("namespace", ns)) +} diff --git a/pkg/buildermgr/namespace_subscriber.go b/pkg/buildermgr/namespace_subscriber.go index f3a9f1d9..cc8e2274 100644 --- a/pkg/buildermgr/namespace_subscriber.go +++ b/pkg/buildermgr/namespace_subscriber.go @@ -11,10 +11,18 @@ type builderEnvNamespaceAdder interface { AddNamespace(ctx context.Context, ns string, mgr manager.Interface) } +type builderEnvNamespaceRemover interface { + RemoveNamespace(ns string) +} + type builderPkgNamespaceAdder interface { AddNamespace(ctx context.Context, ns string, mgr manager.Interface) } +type builderPkgNamespaceRemover interface { + RemoveNamespace(ns string) +} + func NewNamespaceSubscriber(envw builderEnvNamespaceAdder, pkgw builderPkgNamespaceAdder, mgr manager.Interface) utils.NamespaceSubscriber { return utils.NamespaceSubscriberFuncs{ SubscriberName: "buildermgr", @@ -22,6 +30,10 @@ func NewNamespaceSubscriber(envw builderEnvNamespaceAdder, pkgw builderPkgNamesp registerBuilderNamespace(ctx, record.Name, envw, pkgw, mgr) return nil }, + RemoveFunc: func(ctx context.Context, record utils.NamespaceRecord) error { + deregisterBuilderNamespace(record.Name, envw, pkgw) + return nil + }, ResyncFunc: func(ctx context.Context, record utils.NamespaceRecord) error { registerBuilderNamespace(ctx, record.Name, envw, pkgw, mgr) return nil @@ -41,3 +53,16 @@ func registerBuilderNamespace(ctx context.Context, namespace string, envw builde pkgw.AddNamespace(ctx, namespace, mgr) } } + +func deregisterBuilderNamespace(namespace string, envw builderEnvNamespaceAdder, pkgw builderPkgNamespaceAdder) { + if namespace == "" { + return + } + utils.DefaultNSResolver().RemoveNamespace(namespace) + if r, ok := envw.(builderEnvNamespaceRemover); ok { + r.RemoveNamespace(namespace) + } + if r, ok := pkgw.(builderPkgNamespaceRemover); ok { + r.RemoveNamespace(namespace) + } +} diff --git a/pkg/buildermgr/ns_watcher.go b/pkg/buildermgr/ns_watcher.go index 4b2f37f8..ecefdd39 100644 --- a/pkg/buildermgr/ns_watcher.go +++ b/pkg/buildermgr/ns_watcher.go @@ -25,7 +25,9 @@ func StartNSWatcher( pkgw *packageWatcher, mgr manager.Interface, ) { - _, err := utils.RunManagedNamespaceWatcher(ctx, logger, kubeClient, mgr, utils.NewDefaultManagedNamespaceWatcherConfig("buildermgr.NSWatcher", NewNamespaceSubscriber(envw, pkgw, mgr))) + config := utils.NewDefaultManagedNamespaceWatcherConfig("buildermgr.NSWatcher", NewNamespaceSubscriber(envw, pkgw, mgr)) + config.RemovalStrategy = utils.NamespaceRemovalStrategyDispatchRemove + _, err := utils.RunManagedNamespaceWatcher(ctx, logger, kubeClient, mgr, config) if err != nil { logger.Error("buildermgr.NSWatcher: BootstrapAndDispatch failed", zap.Error(err)) } diff --git a/pkg/buildermgr/pkgwatcher.go b/pkg/buildermgr/pkgwatcher.go index 92eac0cc..74bc1bf9 100644 --- a/pkg/buildermgr/pkgwatcher.go +++ b/pkg/buildermgr/pkgwatcher.go @@ -49,6 +49,8 @@ type ( pkgInformer map[string]k8sCache.SharedIndexInformer storageSvcUrl string buildCache *cache.Cache[crd.CacheKeyUR, *fv1.Package] + // nsCancels holds per-namespace context cancel functions. + nsCancels map[string]context.CancelFunc } ) @@ -64,6 +66,7 @@ func makePackageWatcher(logger *zap.Logger, fissionClient versioned.Interface, k pkgInformer: pkgInformer, storageSvcUrl: storageSvcUrl, buildCache: cache.MakeCache[crd.CacheKeyUR, *fv1.Package](0, 0), + nsCancels: make(map[string]context.CancelFunc), } return pkgw } @@ -363,7 +366,23 @@ func (pkgw *packageWatcher) AddNamespace(ctx context.Context, ns string, mgr man ns + "/pkg": pkgInf, ns + "/pod": podInf, }) - fissionFactory.Start(ctx.Done()) - podFactory.Start(ctx.Done()) + + // Create a per-namespace cancellable context for informer lifecycle. + nsCtx, nsCancel := context.WithCancel(ctx) + pkgw.nsCancels[ns] = nsCancel + + fissionFactory.Start(nsCtx.Done()) + podFactory.Start(nsCtx.Done()) pkgw.logger.Info("buildermgr.pkgWatcher.AddNamespace: done", zap.String("namespace", ns)) } + +// RemoveNamespace deregisters a namespace from the package watcher. +func (pkgw *packageWatcher) RemoveNamespace(ns string) { + if cancel, ok := pkgw.nsCancels[ns]; ok { + cancel() + delete(pkgw.nsCancels, ns) + } + delete(pkgw.pkgInformer, ns) + delete(pkgw.podInformer, ns) + pkgw.logger.Info("buildermgr.pkgWatcher.RemoveNamespace: cleaned up", zap.String("namespace", ns)) +} diff --git a/pkg/executor/executortype/container/containermgr.go b/pkg/executor/executortype/container/containermgr.go index 70acec7a..a5082dc4 100644 --- a/pkg/executor/executortype/container/containermgr.go +++ b/pkg/executor/executortype/container/containermgr.go @@ -90,6 +90,10 @@ type ( objectReaperIntervalSecond time.Duration enableOwnerReferences bool + + // nsCancels holds per-namespace context cancel functions so informer + // factories started in AddNamespace can be stopped on RemoveNamespace. + nsCancels map[string]context.CancelFunc } ) @@ -132,8 +136,7 @@ func MakeContainer( deplLister: make(map[string]appslisters.DeploymentLister), deplListerSynced: make(map[string]k8sCache.InformerSynced), svcLister: make(map[string]corelisters.ServiceLister), - svcListerSynced: make(map[string]k8sCache.InformerSynced), - + svcListerSynced: make(map[string]k8sCache.InformerSynced), nsCancels: make(map[string]context.CancelFunc), enableOwnerReferences: utils.IsOwnerReferencesEnabled(), } @@ -833,9 +836,36 @@ func (caaf *Container) AddNamespace(ctx context.Context, ns string, mgr manager. return fmt.Errorf("AddNamespace %s (container): add function handler: %w", ns, err) } - finformer.Start(ctx.Done()) - cnmInformer.Start(ctx.Done()) + // Create a per-namespace cancellable context so RemoveNamespace can stop + // these specific informer factories without affecting the whole process. + nsCtx, nsCancel := context.WithCancel(ctx) + caaf.nsCancels[ns] = nsCancel + + finformer.Start(nsCtx.Done()) + cnmInformer.Start(nsCtx.Done()) caaf.logger.Info("AddNamespace: done (container)", zap.String("namespace", ns)) return nil } + +// RemoveNamespace deregisters a namespace from the container executor. +// Cancels the per-namespace informer context and clears all lister maps so that +// a subsequent AddNamespace call will re-register the namespace correctly. +func (caaf *Container) RemoveNamespace(ctx context.Context, ns string) error { + if ns == "" { + return nil + } + caaf.logger.Info("RemoveNamespace: cleaning up namespace (container)", zap.String("namespace", ns)) + + if cancel, ok := caaf.nsCancels[ns]; ok { + cancel() + delete(caaf.nsCancels, ns) + } + + delete(caaf.deplLister, ns) + delete(caaf.deplListerSynced, ns) + delete(caaf.svcLister, ns) + delete(caaf.svcListerSynced, ns) + + return nil +} diff --git a/pkg/executor/executortype/executortype.go b/pkg/executor/executortype/executortype.go index 6182efee..a36fcf6e 100644 --- a/pkg/executor/executortype/executortype.go +++ b/pkg/executor/executortype/executortype.go @@ -74,4 +74,9 @@ type ExecutorType interface { // starts watching Fission CRDs and K8s resources in it without a pod restart. // Called when a Namespace with label fission.io/managed=true appears. AddNamespace(ctx context.Context, ns string, mgr manager.Interface) error + + // RemoveNamespace deregisters a namespace from the executor, cancelling its + // informer goroutines and clearing dedup state so that a re-add works correctly. + // Called when a Namespace with label fission.io/managed=true is removed. + RemoveNamespace(ctx context.Context, ns string) error } diff --git a/pkg/executor/executortype/newdeploy/newdeploymgr.go b/pkg/executor/executortype/newdeploy/newdeploymgr.go index 1fa34b7d..142ae8b5 100644 --- a/pkg/executor/executortype/newdeploy/newdeploymgr.go +++ b/pkg/executor/executortype/newdeploy/newdeploymgr.go @@ -94,6 +94,10 @@ type ( objectReaperIntervalSecond time.Duration enableOwnerReferences bool + + // nsCancels holds per-namespace context cancel functions so informer + // factories started in AddNamespace can be stopped on RemoveNamespace. + nsCancels map[string]context.CancelFunc } ) @@ -140,8 +144,7 @@ func MakeNewDeploy( deplLister: make(map[string]appslisters.DeploymentLister), deplListerSynced: make(map[string]k8sCache.InformerSynced), svcLister: make(map[string]corelisters.ServiceLister), - svcListerSynced: make(map[string]k8sCache.InformerSynced), - + svcListerSynced: make(map[string]k8sCache.InformerSynced), nsCancels: make(map[string]context.CancelFunc), enableOwnerReferences: utils.IsOwnerReferencesEnabled(), } @@ -949,9 +952,36 @@ func (deploy *NewDeploy) AddNamespace(ctx context.Context, ns string, mgr manage return fmt.Errorf("AddNamespace %s (newdeploy): add environment handler: %w", ns, err) } - finformer.Start(ctx.Done()) - ndmInformer.Start(ctx.Done()) + // Create a per-namespace cancellable context so RemoveNamespace can stop + // these specific informer factories without affecting the whole process. + nsCtx, nsCancel := context.WithCancel(ctx) + deploy.nsCancels[ns] = nsCancel + + finformer.Start(nsCtx.Done()) + ndmInformer.Start(nsCtx.Done()) deploy.logger.Info("AddNamespace: done (newdeploy)", zap.String("namespace", ns)) return nil } + +// RemoveNamespace deregisters a namespace from the newdeploy executor. +// Cancels the per-namespace informer context and clears all lister maps so that +// a subsequent AddNamespace call will re-register the namespace correctly. +func (deploy *NewDeploy) RemoveNamespace(ctx context.Context, ns string) error { + if ns == "" { + return nil + } + deploy.logger.Info("RemoveNamespace: cleaning up namespace (newdeploy)", zap.String("namespace", ns)) + + if cancel, ok := deploy.nsCancels[ns]; ok { + cancel() + delete(deploy.nsCancels, ns) + } + + delete(deploy.deplLister, ns) + delete(deploy.deplListerSynced, ns) + delete(deploy.svcLister, ns) + delete(deploy.svcListerSynced, ns) + + return nil +} diff --git a/pkg/executor/executortype/poolmgr/gpm.go b/pkg/executor/executortype/poolmgr/gpm.go index dcf98fb8..3394fa42 100644 --- a/pkg/executor/executortype/poolmgr/gpm.go +++ b/pkg/executor/executortype/poolmgr/gpm.go @@ -98,6 +98,10 @@ type ( podSpecPatch *apiv1.PodSpec objectReaperIntervalSecond time.Duration + + // nsCancels holds per-namespace context cancel functions so informer + // factories started in AddNamespace can be stopped on RemoveNamespace. + nsCancels map[string]context.CancelFunc } request struct { requestType @@ -159,6 +163,7 @@ func MakeGenericPoolManager(ctx context.Context, objectReaperIntervalSecond: time.Duration(executorUtils.GetObjectReaperInterval(logger, fv1.ExecutorTypePoolmgr, 5)) * time.Second, podLister: make(map[string]corelisters.PodLister), podListerSynced: make(map[string]k8sCache.InformerSynced), + nsCancels: make(map[string]context.CancelFunc), } for ns, informerFactory := range gpmInformerFactory { gpm.podLister[ns] = informerFactory.Core().V1().Pods().Lister() @@ -833,10 +838,40 @@ func (gpm *GenericPoolManager) AddNamespace(ctx context.Context, ns string, mgr return fmt.Errorf("AddNamespace %s: register informers: %w", ns, err) } + // Create a per-namespace cancellable context so RemoveNamespace can stop + // these specific informer factories without affecting the whole process. + nsCtx, nsCancel := context.WithCancel(ctx) + gpm.nsCancels[ns] = nsCancel + // Start the factories — they will begin syncing immediately. - finformer.Start(ctx.Done()) - gpmInformer.Start(ctx.Done()) + finformer.Start(nsCtx.Done()) + gpmInformer.Start(nsCtx.Done()) gpm.logger.Info("AddNamespace: informers started for namespace", zap.String("namespace", ns)) return nil } + +// RemoveNamespace deregisters a namespace from the poolmgr executor. +// Cancels the per-namespace informer context and clears all lister maps so that +// a subsequent AddNamespace call will re-register the namespace correctly. +func (gpm *GenericPoolManager) RemoveNamespace(ctx context.Context, ns string) error { + if ns == "" { + return nil + } + gpm.logger.Info("RemoveNamespace: cleaning up namespace (poolmgr)", zap.String("namespace", ns)) + + // Stop informer factories for this namespace. + if cancel, ok := gpm.nsCancels[ns]; ok { + cancel() + delete(gpm.nsCancels, ns) + } + + // Clear gpm-level lister maps so dedup passes on next AddNamespace. + delete(gpm.podLister, ns) + delete(gpm.podListerSynced, ns) + + // Clear PoolPodController lister maps. + gpm.poolPodC.RemoveNamespace(ns) + + return nil +} diff --git a/pkg/executor/executortype/poolmgr/poolpodcontroller.go b/pkg/executor/executortype/poolmgr/poolpodcontroller.go index 27983f02..d554befc 100644 --- a/pkg/executor/executortype/poolmgr/poolpodcontroller.go +++ b/pkg/executor/executortype/poolmgr/poolpodcontroller.go @@ -528,3 +528,14 @@ func (p *PoolPodController) AddNamespaceInformers( p.logger.Info("AddNamespaceInformers: registered informers for namespace", zap.String("namespace", ns)) return nil } + +// RemoveNamespace clears all per-namespace lister state in the PoolPodController. +// Called from GenericPoolManager.RemoveNamespace so the dedup check in AddNamespace +// will pass if the namespace is re-added later. +func (p *PoolPodController) RemoveNamespace(ns string) { + delete(p.envLister, ns) + delete(p.envListerSynced, ns) + delete(p.podLister, ns) + delete(p.podListerSynced, ns) + p.logger.Info("PoolPodController.RemoveNamespace: cleared lister state", zap.String("namespace", ns)) +} diff --git a/pkg/executor/multitenant/namespace_subscriber.go b/pkg/executor/multitenant/namespace_subscriber.go index 1c6aba67..087f7612 100644 --- a/pkg/executor/multitenant/namespace_subscriber.go +++ b/pkg/executor/multitenant/namespace_subscriber.go @@ -24,6 +24,9 @@ func NewNamespaceSubscriber( registerNamespace(ctx, logger, kubernetesClient, record.Name, executorTypes, mgr) return nil }, + RemoveFunc: func(ctx context.Context, record utils.NamespaceRecord) error { + return deregisterNamespace(ctx, logger, record.Name, executorTypes) + }, ResyncFunc: func(ctx context.Context, record utils.NamespaceRecord) error { registerNamespace(ctx, logger, kubernetesClient, record.Name, executorTypes, mgr) return nil diff --git a/pkg/executor/multitenant/ns_watcher.go b/pkg/executor/multitenant/ns_watcher.go index 27d0dc92..70e68a82 100644 --- a/pkg/executor/multitenant/ns_watcher.go +++ b/pkg/executor/multitenant/ns_watcher.go @@ -87,7 +87,9 @@ func StartNSWatcher( executorTypes map[fv1.ExecutorType]executortype.ExecutorType, mgr manager.Interface, ) { - _, err := utils.RunManagedNamespaceWatcher(ctx, logger, kubernetesClient, mgr, utils.NewDefaultManagedNamespaceWatcherConfig("multitenant.NSWatcher", NewNamespaceSubscriber(logger, kubernetesClient, executorTypes, mgr))) + config := utils.NewDefaultManagedNamespaceWatcherConfig("multitenant.NSWatcher", NewNamespaceSubscriber(logger, kubernetesClient, executorTypes, mgr)) + config.RemovalStrategy = utils.NamespaceRemovalStrategyDispatchRemove + _, err := utils.RunManagedNamespaceWatcher(ctx, logger, kubernetesClient, mgr, config) if err != nil { logger.Error("multitenant.NSWatcher: BootstrapAndDispatch failed", zap.Error(err)) } @@ -135,3 +137,26 @@ func registerExecutorTypes( } return joinErr } + +// deregisterNamespace calls RemoveNamespace on every executor type for the given namespace. +// Called when a Namespace with label fission.io/managed=true is removed. +func deregisterNamespace( + ctx context.Context, + logger *zap.Logger, + ns string, + executorTypes map[fv1.ExecutorType]executortype.ExecutorType, +) error { + var joinErr error + + for _, et := range executorTypes { + if err := et.RemoveNamespace(ctx, ns); err != nil { + logger.Error("multitenant.NSWatcher: RemoveNamespace failed", + zap.String("namespace", ns), + zap.Error(err), + ) + joinErr = errors.Join(joinErr, err) + } + } + logger.Info("multitenant.NSWatcher: deregistered namespace", zap.String("namespace", ns)) + return joinErr +} diff --git a/pkg/executor/multitenant/ns_watcher_test.go b/pkg/executor/multitenant/ns_watcher_test.go index 4f91643a..39204002 100644 --- a/pkg/executor/multitenant/ns_watcher_test.go +++ b/pkg/executor/multitenant/ns_watcher_test.go @@ -46,6 +46,7 @@ func (f *fakeExecutorType) AddNamespace(ctx context.Context, ns string, mgr mana f.lastNamespace = ns return f.addErr } +func (f *fakeExecutorType) RemoveNamespace(ctx context.Context, ns string) error { return nil } var _ executortype.ExecutorType = (*fakeExecutorType)(nil) diff --git a/pkg/router/httpTriggers.go b/pkg/router/httpTriggers.go index e96a4dcf..e12b55dc 100644 --- a/pkg/router/httpTriggers.go +++ b/pkg/router/httpTriggers.go @@ -50,16 +50,18 @@ type HTTPTriggerSet struct { *functionServiceMap *mutableRouter - logger *zap.Logger - fissionClient versioned.Interface - kubeClient kubernetes.Interface - executor eclient.ClientInterface - resolver *functionReferenceResolver - triggers []fv1.HTTPTrigger - triggerInformer map[string]k8sCache.SharedIndexInformer - functions []fv1.Function - funcInformer map[string]k8sCache.SharedIndexInformer - informerMu sync.RWMutex + logger *zap.Logger + fissionClient versioned.Interface + kubeClient kubernetes.Interface + executor eclient.ClientInterface + resolver *functionReferenceResolver + triggers []fv1.HTTPTrigger + triggerInformer map[string]k8sCache.SharedIndexInformer + functions []fv1.Function + funcInformer map[string]k8sCache.SharedIndexInformer + informerMu sync.RWMutex + // nsCancels holds per-namespace context cancel functions for informer lifecycle. + nsCancels map[string]context.CancelFunc updateRouterRequestChannel chan struct{} tsRoundTripperParams *tsRoundTripperParams isDebugEnv bool @@ -84,6 +86,7 @@ func makeHTTPTriggerSet(logger *zap.Logger, fmap *functionServiceMap, fissionCli svcAddrUpdateThrottler: actionThrottler, unTapServiceTimeout: unTapServiceTimeout, syncDebouncer: debounce.New(time.Millisecond * 20), + nsCancels: make(map[string]context.CancelFunc), } httpTriggerSet.triggerInformer = utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.HttpTriggerResource) httpTriggerSet.funcInformer = utils.GetInformersForNamespaces(fissionClient, time.Minute*30, fv1.FunctionResource) @@ -526,11 +529,39 @@ func (ts *HTTPTriggerSet) AddNamespace(ctx context.Context, ns string, mgr manag ns + "/trigger": triggerInf, ns + "/func": funcInf, }) - factory.Start(ctx.Done()) + + // Create a per-namespace cancellable context for informer lifecycle. + nsCtx, nsCancel := context.WithCancel(ctx) + ts.nsCancels[ns] = nsCancel + + factory.Start(nsCtx.Done()) // Wait for cache to sync before rebuilding the router, so triggers are visible. - k8sCache.WaitForCacheSync(ctx.Done(), triggerInf.HasSynced, funcInf.HasSynced) + k8sCache.WaitForCacheSync(nsCtx.Done(), triggerInf.HasSynced, funcInf.HasSynced) ts.logger.Info("router.AddNamespace: done", zap.String("namespace", ns)) ts.syncTriggers() return nil } + +// RemoveNamespace deregisters a namespace from the router. +// Cancels the per-namespace informer context, removes informers from internal maps, +// and triggers a syncTriggers so stale routes are removed immediately. +func (ts *HTTPTriggerSet) RemoveNamespace(ns string) { + if ns == "" { + return + } + ts.logger.Info("router.RemoveNamespace: cleaning up namespace", zap.String("namespace", ns)) + + if cancel, ok := ts.nsCancels[ns]; ok { + cancel() + delete(ts.nsCancels, ns) + } + + ts.informerMu.Lock() + delete(ts.triggerInformer, ns) + delete(ts.funcInformer, ns) + ts.informerMu.Unlock() + + ts.syncTriggers() + ts.logger.Info("router.RemoveNamespace: done", zap.String("namespace", ns)) +} diff --git a/pkg/router/namespace_subscriber.go b/pkg/router/namespace_subscriber.go index 3aedbae0..03195262 100644 --- a/pkg/router/namespace_subscriber.go +++ b/pkg/router/namespace_subscriber.go @@ -11,12 +11,22 @@ type routerNamespaceAdder interface { AddNamespace(ctx context.Context, ns string, mgr manager.Interface) error } +type routerNamespaceRemover interface { + RemoveNamespace(ns string) +} + func NewNamespaceSubscriber(ts routerNamespaceAdder, mgr manager.Interface) utils.NamespaceSubscriber { return utils.NamespaceSubscriberFuncs{ SubscriberName: "router", AddFunc: func(ctx context.Context, record utils.NamespaceRecord) error { return registerRouterNamespace(ctx, record.Name, ts, mgr) }, + RemoveFunc: func(ctx context.Context, record utils.NamespaceRecord) error { + if r, ok := ts.(routerNamespaceRemover); ok { + r.RemoveNamespace(record.Name) + } + return nil + }, ResyncFunc: func(ctx context.Context, record utils.NamespaceRecord) error { return registerRouterNamespace(ctx, record.Name, ts, mgr) }, diff --git a/pkg/router/ns_watcher.go b/pkg/router/ns_watcher.go index 69c0a33f..a23725dd 100644 --- a/pkg/router/ns_watcher.go +++ b/pkg/router/ns_watcher.go @@ -25,7 +25,9 @@ func StartNSWatcher( ts *HTTPTriggerSet, mgr manager.Interface, ) { - _, err := utils.RunManagedNamespaceWatcher(ctx, logger, kubeClient, mgr, utils.NewDefaultManagedNamespaceWatcherConfig("router.NSWatcher", NewNamespaceSubscriber(ts, mgr))) + config := utils.NewDefaultManagedNamespaceWatcherConfig("router.NSWatcher", NewNamespaceSubscriber(ts, mgr)) + config.RemovalStrategy = utils.NamespaceRemovalStrategyDispatchRemove + _, err := utils.RunManagedNamespaceWatcher(ctx, logger, kubeClient, mgr, config) if err != nil { logger.Error("router.NSWatcher: BootstrapAndDispatch failed", zap.Error(err)) }