layer1: guard router informer maps step 6

This commit is contained in:
Naeel
2026-04-26 09:37:09 +03:00
parent 94f26b69ee
commit 87477d4529
3 changed files with 76 additions and 6 deletions
@@ -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.
+11
View File
@@ -18,6 +18,7 @@ package router
import ( import (
"fmt" "fmt"
"sync"
"time" "time"
"go.uber.org/zap" "go.uber.org/zap"
@@ -35,6 +36,7 @@ type (
// FunctionReference -> function metadata // FunctionReference -> function metadata
refCache *cache.Cache[namespacedTriggerReference, resolveResult] refCache *cache.Cache[namespacedTriggerReference, resolveResult]
funcInformer map[string]k8sCache.SharedIndexInformer funcInformer map[string]k8sCache.SharedIndexInformer
mu sync.RWMutex
logger *zap.Logger logger *zap.Logger
// store k8sCache.Store // store k8sCache.Store
} }
@@ -120,12 +122,21 @@ func (frr *functionReferenceResolver) resolve(trigger fv1.HTTPTrigger) (*resolve
} }
func (frr *functionReferenceResolver) getInformerByNamespace(namespace string) (k8sCache.SharedIndexInformer, error) { func (frr *functionReferenceResolver) getInformerByNamespace(namespace string) (k8sCache.SharedIndexInformer, error) {
frr.mu.RLock()
defer frr.mu.RUnlock()
if informer, ok := frr.funcInformer[namespace]; ok { if informer, ok := frr.funcInformer[namespace]; ok {
return informer, nil return informer, nil
} }
return nil, fmt.Errorf("informer for namespace %s not found", namespace) 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. // resolveByName simply looks up function by name in a namespace.
func (frr *functionReferenceResolver) resolveByName(namespace, name string) (*resolveResult, error) { func (frr *functionReferenceResolver) resolveByName(namespace, name string) (*resolveResult, error) {
// get function from cache // get function from cache
+31 -6
View File
@@ -21,6 +21,7 @@ import (
"fmt" "fmt"
"net/http" "net/http"
"strings" "strings"
"sync"
"time" "time"
"github.com/bep/debounce" "github.com/bep/debounce"
@@ -58,6 +59,7 @@ type HTTPTriggerSet struct {
triggerInformer map[string]k8sCache.SharedIndexInformer triggerInformer map[string]k8sCache.SharedIndexInformer
functions []fv1.Function functions []fv1.Function
funcInformer map[string]k8sCache.SharedIndexInformer funcInformer map[string]k8sCache.SharedIndexInformer
informerMu sync.RWMutex
updateRouterRequestChannel chan struct{} updateRouterRequestChannel chan struct{}
tsRoundTripperParams *tsRoundTripperParams tsRoundTripperParams *tsRoundTripperParams
isDebugEnv bool isDebugEnv bool
@@ -310,7 +312,7 @@ func (ts *HTTPTriggerSet) updateTriggerStatusFailed(ht *fv1.HTTPTrigger, err err
} }
func (ts *HTTPTriggerSet) addTriggerHandlers() error { func (ts *HTTPTriggerSet) addTriggerHandlers() error {
for _, triggerInformer := range ts.triggerInformer { for _, triggerInformer := range ts.snapshotTriggerInformers() {
_, err := triggerInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ _, err := triggerInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) { AddFunc: func(obj interface{}) {
trigger := obj.(*fv1.HTTPTrigger) trigger := obj.(*fv1.HTTPTrigger)
@@ -342,7 +344,7 @@ func (ts *HTTPTriggerSet) addTriggerHandlers() error {
} }
func (ts *HTTPTriggerSet) addFunctionHandlers() error { func (ts *HTTPTriggerSet) addFunctionHandlers() error {
for _, funcInformer := range ts.funcInformer { for _, funcInformer := range ts.snapshotFuncInformers() {
_, err := funcInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ _, err := funcInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) { 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) { func (ts *HTTPTriggerSet) updateRouter(ctx context.Context) {
for { for {
select { select {
@@ -398,7 +422,7 @@ func (ts *HTTPTriggerSet) updateRouter(ctx context.Context) {
} }
// get triggers // get triggers
alltriggers := make([]fv1.HTTPTrigger, 0) alltriggers := make([]fv1.HTTPTrigger, 0)
for _, triggerInformer := range ts.triggerInformer { for _, triggerInformer := range ts.snapshotTriggerInformers() {
latestTriggers := triggerInformer.GetStore().List() latestTriggers := triggerInformer.GetStore().List()
for _, t := range latestTriggers { for _, t := range latestTriggers {
alltriggers = append(alltriggers, *t.(*fv1.HTTPTrigger)) alltriggers = append(alltriggers, *t.(*fv1.HTTPTrigger))
@@ -409,7 +433,7 @@ func (ts *HTTPTriggerSet) updateRouter(ctx context.Context) {
// get functions // get functions
allfunctions := make([]fv1.Function, 0) allfunctions := make([]fv1.Function, 0)
functionTimeout := make(map[types.UID]int, 0) functionTimeout := make(map[types.UID]int, 0)
for _, funcInformer := range ts.funcInformer { for _, funcInformer := range ts.snapshotFuncInformers() {
latestFunctions := funcInformer.GetStore().List() latestFunctions := funcInformer.GetStore().List()
for _, f := range latestFunctions { for _, f := range latestFunctions {
fn := *f.(*fv1.Function) 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) return fmt.Errorf("router.AddNamespace %s: func handler: %w", ns, err)
} }
// ts.funcInformer and resolver.funcInformer are the same map reference — ts.informerMu.Lock()
// updating ts.funcInformer also makes the resolver aware of the new namespace.
ts.triggerInformer[ns] = triggerInf ts.triggerInformer[ns] = triggerInf
ts.funcInformer[ns] = funcInf ts.funcInformer[ns] = funcInf
ts.informerMu.Unlock()
ts.resolver.addInformer(ns, funcInf)
mgr.AddInformers(ctx, map[string]k8sCache.SharedIndexInformer{ mgr.AddInformers(ctx, map[string]k8sCache.SharedIndexInformer{
ns + "/trigger": triggerInf, ns + "/trigger": triggerInf,