diff --git a/doc/thinking/2026-04-26-namespace-manager-step38.md b/doc/thinking/2026-04-26-namespace-manager-step38.md new file mode 100644 index 00000000..cce4265a --- /dev/null +++ b/doc/thinking/2026-04-26-namespace-manager-step38.md @@ -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. \ No newline at end of file diff --git a/pkg/buildermgr/ns_watcher.go b/pkg/buildermgr/ns_watcher.go index a55ab8f2..2e197c04 100644 --- a/pkg/buildermgr/ns_watcher.go +++ b/pkg/buildermgr/ns_watcher.go @@ -51,8 +51,7 @@ func StartNSWatcher( if !ok || nsObj.Name == "" { return } - nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventAdd, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())) - if _, _, err := nsManager.DispatchAdd(ctx, nsObj.Name); err != nil { + if _, _, err := utils.DispatchNamespaceAdd(ctx, nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()); err != nil { logger.Error("buildermgr.NSWatcher: DispatchAdd failed", zap.String("namespace", nsObj.Name), zap.Error(err)) } @@ -63,9 +62,8 @@ func StartNSWatcher( return } oldNSObj, _ := oldObj.(*corev1.Namespace) - if oldNSObj != nil && utils.IsManagedNamespace(oldNSObj.Labels) && !utils.IsManagedNamespace(nsObj.Labels) { - event := utils.NamespaceEventFromNamespace(utils.NamespaceEventRemove, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()) - nsManager.Upsert(event) + if utils.NamespaceBecameUnmanaged(oldNSObj, nsObj) { + _, _ = utils.RecordNamespaceRemoval(nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()) logger.Info("buildermgr.NSWatcher: namespace removed from manager state; runtime registrations kept", zap.String("namespace", nsObj.Name)) return @@ -73,20 +71,18 @@ func StartNSWatcher( if !utils.IsManagedNamespace(nsObj.Labels) { return } - nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventUpdate, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())) - if _, _, err := nsManager.DispatchResync(ctx, nsObj.Name); err != nil { + if _, _, err := utils.DispatchNamespaceResync(ctx, nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()); err != nil { logger.Error("buildermgr.NSWatcher: DispatchResync failed", zap.String("namespace", nsObj.Name), zap.Error(err)) } }, DeleteFunc: func(obj interface{}) { - event := utils.NamespaceEventFromObject(utils.NamespaceEventRemove, obj, utils.NamespaceSourceWatcher, time.Now().UTC()) - if event.Name == "" { + record, ok := utils.RecordNamespaceRemoval(nsManager, obj, utils.NamespaceSourceWatcher, time.Now().UTC()) + if !ok { return } - nsManager.Upsert(event) logger.Info("buildermgr.NSWatcher: namespace deleted from manager state; runtime registrations kept", - zap.String("namespace", event.Name)) + zap.String("namespace", record.Name)) }, }) diff --git a/pkg/executor/multitenant/ns_watcher.go b/pkg/executor/multitenant/ns_watcher.go index 2eccfe1f..c45b07f8 100644 --- a/pkg/executor/multitenant/ns_watcher.go +++ b/pkg/executor/multitenant/ns_watcher.go @@ -117,8 +117,7 @@ func StartNSWatcher( if !ok || nsObj.Name == "" { return } - nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventAdd, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())) - if _, _, err := nsManager.DispatchAdd(ctx, nsObj.Name); err != nil { + if _, _, err := utils.DispatchNamespaceAdd(ctx, nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()); err != nil { logger.Error("multitenant.NSWatcher: DispatchAdd failed", zap.String("namespace", nsObj.Name), zap.Error(err)) } @@ -131,9 +130,8 @@ func StartNSWatcher( return } oldNSObj, _ := oldObj.(*corev1.Namespace) - if oldNSObj != nil && utils.IsManagedNamespace(oldNSObj.Labels) && !utils.IsManagedNamespace(nsObj.Labels) { - event := utils.NamespaceEventFromNamespace(utils.NamespaceEventRemove, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()) - nsManager.Upsert(event) + if utils.NamespaceBecameUnmanaged(oldNSObj, nsObj) { + _, _ = utils.RecordNamespaceRemoval(nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()) logger.Info("multitenant.NSWatcher: namespace removed from manager state; runtime registrations kept", zap.String("namespace", nsObj.Name)) return @@ -141,20 +139,18 @@ func StartNSWatcher( if !utils.IsManagedNamespace(nsObj.Labels) { 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 := nsManager.DispatchResync(ctx, nsObj.Name); err != nil { + if _, _, err := utils.DispatchNamespaceResync(ctx, nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()); err != nil { logger.Error("multitenant.NSWatcher: DispatchResync failed", zap.String("namespace", nsObj.Name), zap.Error(err)) } }, DeleteFunc: func(obj interface{}) { - event := utils.NamespaceEventFromObject(utils.NamespaceEventRemove, obj, utils.NamespaceSourceWatcher, time.Now().UTC()) - if event.Name == "" { + record, ok := utils.RecordNamespaceRemoval(nsManager, obj, utils.NamespaceSourceWatcher, time.Now().UTC()) + if !ok { return } - nsManager.Upsert(event) logger.Info("multitenant.NSWatcher: namespace deleted from manager state; runtime registrations kept", - zap.String("namespace", event.Name)) + zap.String("namespace", record.Name)) }, }) diff --git a/pkg/router/ns_watcher.go b/pkg/router/ns_watcher.go index 787357a3..dab55ef6 100644 --- a/pkg/router/ns_watcher.go +++ b/pkg/router/ns_watcher.go @@ -51,8 +51,7 @@ func StartNSWatcher( if !ok || nsObj.Name == "" { return } - nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventAdd, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())) - if _, _, err := nsManager.DispatchAdd(ctx, nsObj.Name); err != nil { + if _, _, err := utils.DispatchNamespaceAdd(ctx, nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()); err != nil { logger.Error("router.NSWatcher: DispatchAdd failed", zap.String("namespace", nsObj.Name), zap.Error(err)) } @@ -63,9 +62,8 @@ func StartNSWatcher( return } oldNSObj, _ := oldObj.(*corev1.Namespace) - if oldNSObj != nil && utils.IsManagedNamespace(oldNSObj.Labels) && !utils.IsManagedNamespace(nsObj.Labels) { - event := utils.NamespaceEventFromNamespace(utils.NamespaceEventRemove, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()) - nsManager.Upsert(event) + if utils.NamespaceBecameUnmanaged(oldNSObj, nsObj) { + _, _ = utils.RecordNamespaceRemoval(nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()) logger.Info("router.NSWatcher: namespace removed from manager state; runtime registrations kept", zap.String("namespace", nsObj.Name)) return @@ -73,20 +71,18 @@ func StartNSWatcher( if !utils.IsManagedNamespace(nsObj.Labels) { return } - nsManager.Upsert(utils.NamespaceEventFromNamespace(utils.NamespaceEventUpdate, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC())) - if _, _, err := nsManager.DispatchResync(ctx, nsObj.Name); err != nil { + if _, _, err := utils.DispatchNamespaceResync(ctx, nsManager, nsObj, utils.NamespaceSourceWatcher, time.Now().UTC()); err != nil { logger.Error("router.NSWatcher: DispatchResync failed", zap.String("namespace", nsObj.Name), zap.Error(err)) } }, DeleteFunc: func(obj interface{}) { - event := utils.NamespaceEventFromObject(utils.NamespaceEventRemove, obj, utils.NamespaceSourceWatcher, time.Now().UTC()) - if event.Name == "" { + record, ok := utils.RecordNamespaceRemoval(nsManager, obj, utils.NamespaceSourceWatcher, time.Now().UTC()) + if !ok { return } - nsManager.Upsert(event) logger.Info("router.NSWatcher: namespace deleted from manager state; runtime registrations kept", - zap.String("namespace", event.Name)) + zap.String("namespace", record.Name)) }, }) diff --git a/pkg/utils/namespace_manager.go b/pkg/utils/namespace_manager.go index 4fa70ebb..9e92974b 100644 --- a/pkg/utils/namespace_manager.go +++ b/pkg/utils/namespace_manager.go @@ -6,6 +6,8 @@ import ( "sort" "sync" "time" + + corev1 "k8s.io/api/core/v1" ) type NamespaceSubscriber interface { @@ -100,6 +102,41 @@ func NewWatcherNamespaceManager(ctx context.Context, namespaces []string, source 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) { m.mu.Lock() defer m.mu.Unlock() diff --git a/pkg/utils/namespace_manager_test.go b/pkg/utils/namespace_manager_test.go index 975ffec5..cff69d4a 100644 --- a/pkg/utils/namespace_manager_test.go +++ b/pkg/utils/namespace_manager_test.go @@ -6,6 +6,9 @@ import ( "reflect" "testing" "time" + + corev1 "k8s.io/api/core/v1" + k8sCache "k8s.io/client-go/tools/cache" ) 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) { manager := NewNamespaceManager() manager.Upsert(NamespaceEvent{Type: NamespaceEventAdd, Name: "tenant-a", Source: NamespaceSourceWatcher})