From 87477d4529c8a3d2b67de9f6bf3c90bbb846d9dd Mon Sep 17 00:00:00 2001 From: Naeel Date: Sun, 26 Apr 2026 09:37:09 +0300 Subject: [PATCH] layer1: guard router informer maps step 6 --- .../2026-04-26-namespace-manager-step6.md | 34 +++++++++++++++++ pkg/router/functionReferenceResolver.go | 11 ++++++ pkg/router/httpTriggers.go | 37 ++++++++++++++++--- 3 files changed, 76 insertions(+), 6 deletions(-) create mode 100644 doc/thinking/2026-04-26-namespace-manager-step6.md diff --git a/doc/thinking/2026-04-26-namespace-manager-step6.md b/doc/thinking/2026-04-26-namespace-manager-step6.md new file mode 100644 index 00000000..d304d0d6 --- /dev/null +++ b/doc/thinking/2026-04-26-namespace-manager-step6.md @@ -0,0 +1,34 @@ +# 2026-04-26 — NamespaceManager rewrite, step 6 + +## Цель шага + +Закрыть race-surface в router вокруг динамического добавления namespace informer-ов. + +## Проблема + +В router есть два связанных mutable map: + +- `HTTPTriggerSet.triggerInformer` +- `HTTPTriggerSet.funcInformer` + +`AddNamespace()` пишет в них на лету, а `updateRouter()` одновременно итерируется по ним. +Кроме того, `functionReferenceResolver` получает `funcInformer` map и читает ее без синхронизации. + +Это делает dynamic onboarding потенциальным источником: + +- `concurrent map iteration and map write`; +- чтения неполного снимка informer-ов; +- гонок между router rebuild и resolver lookup. + +## Исправление + +1. В `HTTPTriggerSet` добавляется `RWMutex` для informer maps. +2. Чтение informer-ов переводится на snapshot helpers. +3. `functionReferenceResolver` получает собственный lock и метод `addInformer()`. +4. `router.AddNamespace()` обновляет router map и resolver map под контролируемым доступом. + +## Что НЕ меняем + +- не переписываем router lifecycle целиком; +- не добавляем remove semantics; +- не меняем trigger/function business logic. \ No newline at end of file diff --git a/pkg/router/functionReferenceResolver.go b/pkg/router/functionReferenceResolver.go index 9a03e34e..6ebb86b3 100644 --- a/pkg/router/functionReferenceResolver.go +++ b/pkg/router/functionReferenceResolver.go @@ -18,6 +18,7 @@ package router import ( "fmt" + "sync" "time" "go.uber.org/zap" @@ -35,6 +36,7 @@ type ( // FunctionReference -> function metadata refCache *cache.Cache[namespacedTriggerReference, resolveResult] funcInformer map[string]k8sCache.SharedIndexInformer + mu sync.RWMutex logger *zap.Logger // store k8sCache.Store } @@ -120,12 +122,21 @@ func (frr *functionReferenceResolver) resolve(trigger fv1.HTTPTrigger) (*resolve } func (frr *functionReferenceResolver) getInformerByNamespace(namespace string) (k8sCache.SharedIndexInformer, error) { + frr.mu.RLock() + defer frr.mu.RUnlock() + if informer, ok := frr.funcInformer[namespace]; ok { return informer, nil } return nil, fmt.Errorf("informer for namespace %s not found", namespace) } +func (frr *functionReferenceResolver) addInformer(namespace string, informer k8sCache.SharedIndexInformer) { + frr.mu.Lock() + defer frr.mu.Unlock() + frr.funcInformer[namespace] = informer +} + // resolveByName simply looks up function by name in a namespace. func (frr *functionReferenceResolver) resolveByName(namespace, name string) (*resolveResult, error) { // get function from cache diff --git a/pkg/router/httpTriggers.go b/pkg/router/httpTriggers.go index 26a08c61..e96a4dcf 100644 --- a/pkg/router/httpTriggers.go +++ b/pkg/router/httpTriggers.go @@ -21,6 +21,7 @@ import ( "fmt" "net/http" "strings" + "sync" "time" "github.com/bep/debounce" @@ -58,6 +59,7 @@ type HTTPTriggerSet struct { triggerInformer map[string]k8sCache.SharedIndexInformer functions []fv1.Function funcInformer map[string]k8sCache.SharedIndexInformer + informerMu sync.RWMutex updateRouterRequestChannel chan struct{} tsRoundTripperParams *tsRoundTripperParams isDebugEnv bool @@ -310,7 +312,7 @@ func (ts *HTTPTriggerSet) updateTriggerStatusFailed(ht *fv1.HTTPTrigger, err err } func (ts *HTTPTriggerSet) addTriggerHandlers() error { - for _, triggerInformer := range ts.triggerInformer { + for _, triggerInformer := range ts.snapshotTriggerInformers() { _, err := triggerInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { trigger := obj.(*fv1.HTTPTrigger) @@ -342,7 +344,7 @@ func (ts *HTTPTriggerSet) addTriggerHandlers() error { } func (ts *HTTPTriggerSet) addFunctionHandlers() error { - for _, funcInformer := range ts.funcInformer { + for _, funcInformer := range ts.snapshotFuncInformers() { _, err := funcInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { @@ -389,6 +391,28 @@ func (ts *HTTPTriggerSet) syncTriggers() { }) } +func (ts *HTTPTriggerSet) snapshotTriggerInformers() []k8sCache.SharedIndexInformer { + ts.informerMu.RLock() + defer ts.informerMu.RUnlock() + + informers := make([]k8sCache.SharedIndexInformer, 0, len(ts.triggerInformer)) + for _, informer := range ts.triggerInformer { + informers = append(informers, informer) + } + return informers +} + +func (ts *HTTPTriggerSet) snapshotFuncInformers() []k8sCache.SharedIndexInformer { + ts.informerMu.RLock() + defer ts.informerMu.RUnlock() + + informers := make([]k8sCache.SharedIndexInformer, 0, len(ts.funcInformer)) + for _, informer := range ts.funcInformer { + informers = append(informers, informer) + } + return informers +} + func (ts *HTTPTriggerSet) updateRouter(ctx context.Context) { for { select { @@ -398,7 +422,7 @@ func (ts *HTTPTriggerSet) updateRouter(ctx context.Context) { } // get triggers alltriggers := make([]fv1.HTTPTrigger, 0) - for _, triggerInformer := range ts.triggerInformer { + for _, triggerInformer := range ts.snapshotTriggerInformers() { latestTriggers := triggerInformer.GetStore().List() for _, t := range latestTriggers { alltriggers = append(alltriggers, *t.(*fv1.HTTPTrigger)) @@ -409,7 +433,7 @@ func (ts *HTTPTriggerSet) updateRouter(ctx context.Context) { // get functions allfunctions := make([]fv1.Function, 0) functionTimeout := make(map[types.UID]int, 0) - for _, funcInformer := range ts.funcInformer { + for _, funcInformer := range ts.snapshotFuncInformers() { latestFunctions := funcInformer.GetStore().List() for _, f := range latestFunctions { fn := *f.(*fv1.Function) @@ -492,10 +516,11 @@ func (ts *HTTPTriggerSet) AddNamespace(ctx context.Context, ns string, mgr manag return fmt.Errorf("router.AddNamespace %s: func handler: %w", ns, err) } - // ts.funcInformer and resolver.funcInformer are the same map reference — - // updating ts.funcInformer also makes the resolver aware of the new namespace. + ts.informerMu.Lock() ts.triggerInformer[ns] = triggerInf ts.funcInformer[ns] = funcInf + ts.informerMu.Unlock() + ts.resolver.addInformer(ns, funcInf) mgr.AddInformers(ctx, map[string]k8sCache.SharedIndexInformer{ ns + "/trigger": triggerInf,