layer1: share namespace watcher lifecycle helpers

This commit is contained in:
Naeel
2026-04-26 10:45:03 +03:00
parent 0e08664ef6
commit 4cd4bc9507
6 changed files with 129 additions and 33 deletions
@@ -0,0 +1,19 @@
# 2026-04-26 — NamespaceManager rewrite, step 38
## Цель шага
Убрать повторяющуюся lifecycle логiku namespace watcher-ов.
## Что меняем
1. В `utils` добавляем helpers:
- `NamespaceBecameUnmanaged()`
- `DispatchNamespaceAdd()`
- `DispatchNamespaceResync()`
- `RecordNamespaceRemoval()`
2. `buildermgr`, `router`, `executor/multitenant` используют эти helpers.
## Что НЕ меняем
- не меняем runtime semantics;
- remove по-прежнему только bookkeeping, без cleanup.
+7 -11
View File
@@ -51,8 +51,7 @@ func StartNSWatcher(
if !ok || nsObj.Name == "" { if !ok || nsObj.Name == "" {
return return
} }
nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventAdd, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())) if _, _, err := utils.DispatchNamespaceAdd(ctx, nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()); err != nil {
if _, _, err := nsManager.DispatchAdd(ctx, nsObj.Name); err != nil {
logger.Error("buildermgr.NSWatcher: DispatchAdd failed", logger.Error("buildermgr.NSWatcher: DispatchAdd failed",
zap.String("namespace", nsObj.Name), zap.Error(err)) zap.String("namespace", nsObj.Name), zap.Error(err))
} }
@@ -63,9 +62,8 @@ func StartNSWatcher(
return return
} }
oldNSObj, _ := oldObj.(*corev1.Namespace) oldNSObj, _ := oldObj.(*corev1.Namespace)
if oldNSObj != nil && utils.IsManagedNamespace(oldNSObj.Labels) && !utils.IsManagedNamespace(nsObj.Labels) { if utils.NamespaceBecameUnmanaged(oldNSObj, nsObj) {
event := utils.NamespaceEventFromNamespace(utils.NamespaceEventRemove, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()) _, _ = utils.RecordNamespaceRemoval(nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())
nsManager.Upsert(event)
logger.Info("buildermgr.NSWatcher: namespace removed from manager state; runtime registrations kept", logger.Info("buildermgr.NSWatcher: namespace removed from manager state; runtime registrations kept",
zap.String("namespace", nsObj.Name)) zap.String("namespace", nsObj.Name))
return return
@@ -73,20 +71,18 @@ func StartNSWatcher(
if !utils.IsManagedNamespace(nsObj.Labels) { if !utils.IsManagedNamespace(nsObj.Labels) {
return return
} }
nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventUpdate, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())) if _, _, err := utils.DispatchNamespaceResync(ctx, nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()); err != nil {
if _, _, err := nsManager.DispatchResync(ctx, nsObj.Name); err != nil {
logger.Error("buildermgr.NSWatcher: DispatchResync failed", logger.Error("buildermgr.NSWatcher: DispatchResync failed",
zap.String("namespace", nsObj.Name), zap.Error(err)) zap.String("namespace", nsObj.Name), zap.Error(err))
} }
}, },
DeleteFunc: func(obj interface{}) { DeleteFunc: func(obj interface{}) {
event := utils.NamespaceEventFromObject(utils.NamespaceEventRemove, obj, utils.NamespaceSourceWatcher, time.Now().UTC()) record, ok := utils.RecordNamespaceRemoval(nsManager, obj, utils.NamespaceSourceWatcher, time.Now().UTC())
if event.Name == "" { if !ok {
return return
} }
nsManager.Upsert(event)
logger.Info("buildermgr.NSWatcher: namespace deleted from manager state; runtime registrations kept", logger.Info("buildermgr.NSWatcher: namespace deleted from manager state; runtime registrations kept",
zap.String("namespace", event.Name)) zap.String("namespace", record.Name))
}, },
}) })
+7 -11
View File
@@ -117,8 +117,7 @@ func StartNSWatcher(
if !ok || nsObj.Name == "" { if !ok || nsObj.Name == "" {
return return
} }
nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventAdd, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())) if _, _, err := utils.DispatchNamespaceAdd(ctx, nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()); err != nil {
if _, _, err := nsManager.DispatchAdd(ctx, nsObj.Name); err != nil {
logger.Error("multitenant.NSWatcher: DispatchAdd failed", logger.Error("multitenant.NSWatcher: DispatchAdd failed",
zap.String("namespace", nsObj.Name), zap.Error(err)) zap.String("namespace", nsObj.Name), zap.Error(err))
} }
@@ -131,9 +130,8 @@ func StartNSWatcher(
return return
} }
oldNSObj, _ := oldObj.(*corev1.Namespace) oldNSObj, _ := oldObj.(*corev1.Namespace)
if oldNSObj != nil && utils.IsManagedNamespace(oldNSObj.Labels) && !utils.IsManagedNamespace(nsObj.Labels) { if utils.NamespaceBecameUnmanaged(oldNSObj, nsObj) {
event := utils.NamespaceEventFromNamespace(utils.NamespaceEventRemove, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()) _, _ = utils.RecordNamespaceRemoval(nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())
nsManager.Upsert(event)
logger.Info("multitenant.NSWatcher: namespace removed from manager state; runtime registrations kept", logger.Info("multitenant.NSWatcher: namespace removed from manager state; runtime registrations kept",
zap.String("namespace", nsObj.Name)) zap.String("namespace", nsObj.Name))
return return
@@ -141,20 +139,18 @@ func StartNSWatcher(
if !utils.IsManagedNamespace(nsObj.Labels) { if !utils.IsManagedNamespace(nsObj.Labels) {
return // label was removed — nothing to do (executor keeps existing registrations) return // label was removed — nothing to do (executor keeps existing registrations)
} }
nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventUpdate, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())) if _, _, err := utils.DispatchNamespaceResync(ctx, nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()); err != nil {
if _, _, err := nsManager.DispatchResync(ctx, nsObj.Name); err != nil {
logger.Error("multitenant.NSWatcher: DispatchResync failed", logger.Error("multitenant.NSWatcher: DispatchResync failed",
zap.String("namespace", nsObj.Name), zap.Error(err)) zap.String("namespace", nsObj.Name), zap.Error(err))
} }
}, },
DeleteFunc: func(obj interface{}) { DeleteFunc: func(obj interface{}) {
event := utils.NamespaceEventFromObject(utils.NamespaceEventRemove, obj, utils.NamespaceSourceWatcher, time.Now().UTC()) record, ok := utils.RecordNamespaceRemoval(nsManager, obj, utils.NamespaceSourceWatcher, time.Now().UTC())
if event.Name == "" { if !ok {
return return
} }
nsManager.Upsert(event)
logger.Info("multitenant.NSWatcher: namespace deleted from manager state; runtime registrations kept", logger.Info("multitenant.NSWatcher: namespace deleted from manager state; runtime registrations kept",
zap.String("namespace", event.Name)) zap.String("namespace", record.Name))
}, },
}) })
+7 -11
View File
@@ -51,8 +51,7 @@ func StartNSWatcher(
if !ok || nsObj.Name == "" { if !ok || nsObj.Name == "" {
return return
} }
nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventAdd, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())) if _, _, err := utils.DispatchNamespaceAdd(ctx, nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()); err != nil {
if _, _, err := nsManager.DispatchAdd(ctx, nsObj.Name); err != nil {
logger.Error("router.NSWatcher: DispatchAdd failed", logger.Error("router.NSWatcher: DispatchAdd failed",
zap.String("namespace", nsObj.Name), zap.Error(err)) zap.String("namespace", nsObj.Name), zap.Error(err))
} }
@@ -63,9 +62,8 @@ func StartNSWatcher(
return return
} }
oldNSObj, _ := oldObj.(*corev1.Namespace) oldNSObj, _ := oldObj.(*corev1.Namespace)
if oldNSObj != nil && utils.IsManagedNamespace(oldNSObj.Labels) && !utils.IsManagedNamespace(nsObj.Labels) { if utils.NamespaceBecameUnmanaged(oldNSObj, nsObj) {
event := utils.NamespaceEventFromNamespace(utils.NamespaceEventRemove, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()) _, _ = utils.RecordNamespaceRemoval(nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())
nsManager.Upsert(event)
logger.Info("router.NSWatcher: namespace removed from manager state; runtime registrations kept", logger.Info("router.NSWatcher: namespace removed from manager state; runtime registrations kept",
zap.String("namespace", nsObj.Name)) zap.String("namespace", nsObj.Name))
return return
@@ -73,20 +71,18 @@ func StartNSWatcher(
if !utils.IsManagedNamespace(nsObj.Labels) { if !utils.IsManagedNamespace(nsObj.Labels) {
return return
} }
nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventUpdate, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())) if _, _, err := utils.DispatchNamespaceResync(ctx, nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()); err != nil {
if _, _, err := nsManager.DispatchResync(ctx, nsObj.Name); err != nil {
logger.Error("router.NSWatcher: DispatchResync failed", logger.Error("router.NSWatcher: DispatchResync failed",
zap.String("namespace", nsObj.Name), zap.Error(err)) zap.String("namespace", nsObj.Name), zap.Error(err))
} }
}, },
DeleteFunc: func(obj interface{}) { DeleteFunc: func(obj interface{}) {
event := utils.NamespaceEventFromObject(utils.NamespaceEventRemove, obj, utils.NamespaceSourceWatcher, time.Now().UTC()) record, ok := utils.RecordNamespaceRemoval(nsManager, obj, utils.NamespaceSourceWatcher, time.Now().UTC())
if event.Name == "" { if !ok {
return return
} }
nsManager.Upsert(event)
logger.Info("router.NSWatcher: namespace deleted from manager state; runtime registrations kept", logger.Info("router.NSWatcher: namespace deleted from manager state; runtime registrations kept",
zap.String("namespace", event.Name)) zap.String("namespace", record.Name))
}, },
}) })
+37
View File
@@ -6,6 +6,8 @@ import (
"sort" "sort"
"sync" "sync"
"time" "time"
corev1 "k8s.io/api/core/v1"
) )
type NamespaceSubscriber interface { type NamespaceSubscriber interface {
@@ -100,6 +102,41 @@ func NewWatcherNamespaceManager(ctx context.Context, namespaces []string, source
return manager, err return manager, err
} }
func NamespaceBecameUnmanaged(oldNamespace *corev1.Namespace, newNamespace *corev1.Namespace) bool {
if oldNamespace == nil || newNamespace == nil {
return false
}
return IsManagedNamespace(oldNamespace.Labels) && !IsManagedNamespace(newNamespace.Labels)
}
func DispatchNamespaceAdd(ctx context.Context, manager NamespaceManager, namespace *corev1.Namespace, source NamespaceSource, observedAt time.Time) (NamespaceRecord, bool, error) {
if manager == nil || namespace == nil || namespace.Name == "" {
return NamespaceRecord{}, false, nil
}
manager.Upsert(NamespaceEventFromNamespace(NamespaceEventAdd, namespace, source, observedAt))
return manager.DispatchAdd(ctx, namespace.Name)
}
func DispatchNamespaceResync(ctx context.Context, manager NamespaceManager, namespace *corev1.Namespace, source NamespaceSource, observedAt time.Time) (NamespaceRecord, bool, error) {
if manager == nil || namespace == nil || namespace.Name == "" {
return NamespaceRecord{}, false, nil
}
manager.Upsert(NamespaceEventFromNamespace(NamespaceEventUpdate, namespace, source, observedAt))
return manager.DispatchResync(ctx, namespace.Name)
}
func RecordNamespaceRemoval(manager NamespaceManager, obj interface{}, source NamespaceSource, observedAt time.Time) (NamespaceRecord, bool) {
if manager == nil {
return NamespaceRecord{}, false
}
event := NamespaceEventFromObject(NamespaceEventRemove, obj, source, observedAt)
if event.Name == "" {
return NamespaceRecord{}, false
}
record := manager.Upsert(event)
return record, true
}
func (m *inMemoryNamespaceManager) Subscribe(subscriber NamespaceSubscriber) { func (m *inMemoryNamespaceManager) Subscribe(subscriber NamespaceSubscriber) {
m.mu.Lock() m.mu.Lock()
defer m.mu.Unlock() defer m.mu.Unlock()
+52
View File
@@ -6,6 +6,9 @@ import (
"reflect" "reflect"
"testing" "testing"
"time" "time"
corev1 "k8s.io/api/core/v1"
k8sCache "k8s.io/client-go/tools/cache"
) )
type testNamespaceSubscriber struct { type testNamespaceSubscriber struct {
@@ -285,6 +288,55 @@ func TestNewWatcherNamespaceManager(t *testing.T) {
} }
} }
func TestNamespaceBecameUnmanaged(t *testing.T) {
oldNamespace := &corev1.Namespace{}
oldNamespace.Labels = map[string]string{ManagedNamespaceLabelKey: ManagedNamespaceLabelValue}
newNamespace := &corev1.Namespace{}
newNamespace.Labels = map[string]string{}
if !NamespaceBecameUnmanaged(oldNamespace, newNamespace) {
t.Fatalf("expected namespace to become unmanaged")
}
if NamespaceBecameUnmanaged(nil, newNamespace) {
t.Fatalf("expected nil old namespace to report false")
}
}
func TestDispatchNamespaceAddAndResync(t *testing.T) {
manager := NewNamespaceManager()
router := &testNamespaceSubscriber{name: "router"}
manager.Subscribe(router)
namespace := &corev1.Namespace{}
namespace.Name = "tenant-a"
namespace.Labels = map[string]string{ManagedNamespaceLabelKey: ManagedNamespaceLabelValue}
record, ok, err := DispatchNamespaceAdd(context.Background(), manager, namespace, NamespaceSourceWatcher, time.Now().UTC())
if !ok || err != nil || record.Phase != NamespacePhaseActive {
t.Fatalf("expected add dispatch success, ok=%v err=%v phase=%s", ok, err, record.Phase)
}
record, ok, err = DispatchNamespaceResync(context.Background(), manager, namespace, NamespaceSourceWatcher, time.Now().UTC())
if !ok || err != nil || record.Phase != NamespacePhaseActive {
t.Fatalf("expected resync dispatch success, ok=%v err=%v phase=%s", ok, err, record.Phase)
}
if router.addCalls != 1 || router.resyncCalls != 1 {
t.Fatalf("expected add and resync subscriber calls")
}
}
func TestRecordNamespaceRemoval(t *testing.T) {
manager := NewNamespaceManager()
manager.Upsert(NamespaceEvent{Type: NamespaceEventAdd, Name: "tenant-a", Source: NamespaceSourceWatcher})
namespace := &corev1.Namespace{}
namespace.Name = "tenant-a"
record, ok := RecordNamespaceRemoval(manager, k8sCache.DeletedFinalStateUnknown{Obj: namespace}, NamespaceSourceWatcher, time.Now().UTC())
if !ok {
t.Fatalf("expected removal bookkeeping to succeed")
}
if record.Phase != NamespacePhaseRemoved {
t.Fatalf("expected removed phase after bookkeeping, got %s", record.Phase)
}
}
func TestNamespaceManagerDispatchAdd(t *testing.T) { func TestNamespaceManagerDispatchAdd(t *testing.T) {
manager := NewNamespaceManager() manager := NewNamespaceManager()
manager.Upsert(NamespaceEvent{Type: NamespaceEventAdd, Name: "tenant-a", Source: NamespaceSourceWatcher}) manager.Upsert(NamespaceEvent{Type: NamespaceEventAdd, Name: "tenant-a", Source: NamespaceSourceWatcher})