From b1e2e6462dfb199918afb635aafe6cac79106a0c Mon Sep 17 00:00:00 2001 From: Naeel Date: Sun, 26 Apr 2026 10:22:35 +0300 Subject: [PATCH] layer1: add namespace dispatch step 17 --- .../2026-04-26-namespace-manager-step17.md | 22 ++++++ pkg/utils/namespace_manager.go | 56 ++++++++++++++ pkg/utils/namespace_manager_test.go | 73 ++++++++++++++++--- 3 files changed, 139 insertions(+), 12 deletions(-) create mode 100644 doc/thinking/2026-04-26-namespace-manager-step17.md diff --git a/doc/thinking/2026-04-26-namespace-manager-step17.md b/doc/thinking/2026-04-26-namespace-manager-step17.md new file mode 100644 index 00000000..5fbf5e48 --- /dev/null +++ b/doc/thinking/2026-04-26-namespace-manager-step17.md @@ -0,0 +1,22 @@ +# 2026-04-26 — NamespaceManager rewrite, step 17 + +## Цель шага + +Добавить dispatch helper для прогона namespace через subscriber-ов в add/resync path. + +## Что меняем + +1. В manager interface добавляем: + - `DispatchAdd()` + - `DispatchResync()` +2. Manager вызывает subscriber-ов последовательно. +3. Для каждого subscriber-а manager проставляет part state: + - `registering` + - `active` или `failed` +4. Добавляем unit tests на success и failure path. + +## Что НЕ меняем + +- не подключаем dispatch к production watcher-ам; +- не добавляем remove dispatch; +- не меняем runtime components. \ No newline at end of file diff --git a/pkg/utils/namespace_manager.go b/pkg/utils/namespace_manager.go index a41fe15a..705d2e15 100644 --- a/pkg/utils/namespace_manager.go +++ b/pkg/utils/namespace_manager.go @@ -22,6 +22,8 @@ type NamespaceManager interface { Subscribe(subscriber NamespaceSubscriber) SnapshotSubscribers() []string Upsert(event NamespaceEvent) NamespaceRecord + DispatchAdd(ctx context.Context, namespace string) (NamespaceRecord, bool, error) + DispatchResync(ctx context.Context, namespace string) (NamespaceRecord, bool, error) MarkPartState(namespace string, part string, state NamespacePartState) (NamespaceRecord, bool) MarkPartRegistering(namespace string, part string) (NamespaceRecord, bool) MarkPartActive(namespace string, part string) (NamespaceRecord, bool) @@ -85,6 +87,23 @@ func (m *inMemoryNamespaceManager) SnapshotSubscribers() []string { return names } +func (m *inMemoryNamespaceManager) snapshotSubscriberObjects() []NamespaceSubscriber { + m.mu.RLock() + defer m.mu.RUnlock() + + names := make([]string, 0, len(m.subs)) + for name := range m.subs { + names = append(names, name) + } + sort.Strings(names) + + subscribers := make([]NamespaceSubscriber, 0, len(names)) + for _, name := range names { + subscribers = append(subscribers, m.subs[name]) + } + return subscribers +} + func (m *inMemoryNamespaceManager) Snapshot() []string { m.mu.RLock() defer m.mu.RUnlock() @@ -170,6 +189,43 @@ func (m *inMemoryNamespaceManager) Upsert(event NamespaceEvent) NamespaceRecord return record.Clone() } +func (m *inMemoryNamespaceManager) DispatchAdd(ctx context.Context, namespace string) (NamespaceRecord, bool, error) { + return m.dispatch(ctx, namespace, func(subscriber NamespaceSubscriber, record NamespaceRecord) error { + return subscriber.OnNamespaceAdd(ctx, record) + }) +} + +func (m *inMemoryNamespaceManager) DispatchResync(ctx context.Context, namespace string) (NamespaceRecord, bool, error) { + return m.dispatch(ctx, namespace, func(subscriber NamespaceSubscriber, record NamespaceRecord) error { + return subscriber.OnNamespaceResync(ctx, record) + }) +} + +func (m *inMemoryNamespaceManager) dispatch(ctx context.Context, namespace string, handler func(NamespaceSubscriber, NamespaceRecord) error) (NamespaceRecord, bool, error) { + record, ok := m.Get(namespace) + if !ok { + return NamespaceRecord{}, false, nil + } + + var firstErr error + for _, subscriber := range m.snapshotSubscriberObjects() { + _, _ = m.MarkPartRegistering(namespace, subscriber.Name()) + currentRecord, _ := m.Get(namespace) + err := handler(subscriber, currentRecord) + if err != nil { + _, _ = m.MarkPartFailed(namespace, subscriber.Name(), err) + if firstErr == nil { + firstErr = err + } + continue + } + _, _ = m.MarkPartActive(namespace, subscriber.Name()) + } + + record, _ = m.Get(namespace) + return record, true, firstErr +} + func (m *inMemoryNamespaceManager) MarkPartState(namespace string, part string, state NamespacePartState) (NamespaceRecord, bool) { m.mu.Lock() defer m.mu.Unlock() diff --git a/pkg/utils/namespace_manager_test.go b/pkg/utils/namespace_manager_test.go index 76f0dafe..b812874c 100644 --- a/pkg/utils/namespace_manager_test.go +++ b/pkg/utils/namespace_manager_test.go @@ -9,23 +9,29 @@ import ( ) type testNamespaceSubscriber struct { - name string + name string + addErr error + resyncErr error + addCalls int + resyncCalls int } -func (s testNamespaceSubscriber) Name() string { +func (s *testNamespaceSubscriber) Name() string { return s.name } -func (s testNamespaceSubscriber) OnNamespaceAdd(ctx context.Context, record NamespaceRecord) error { +func (s *testNamespaceSubscriber) OnNamespaceAdd(ctx context.Context, record NamespaceRecord) error { + s.addCalls++ + return s.addErr +} + +func (s *testNamespaceSubscriber) OnNamespaceRemove(ctx context.Context, record NamespaceRecord) error { return nil } -func (s testNamespaceSubscriber) OnNamespaceRemove(ctx context.Context, record NamespaceRecord) error { - return nil -} - -func (s testNamespaceSubscriber) OnNamespaceResync(ctx context.Context, record NamespaceRecord) error { - return nil +func (s *testNamespaceSubscriber) OnNamespaceResync(ctx context.Context, record NamespaceRecord) error { + s.resyncCalls++ + return s.resyncErr } func TestNamespaceManagerSnapshotAndGet(t *testing.T) { @@ -122,9 +128,9 @@ func TestNamespaceManagerSnapshotRecordsReturnsCopies(t *testing.T) { func TestNamespaceManagerSubscribers(t *testing.T) { manager := NewNamespaceManager() - manager.Subscribe(testNamespaceSubscriber{name: "router"}) - manager.Subscribe(testNamespaceSubscriber{name: "buildermgr"}) - manager.Subscribe(testNamespaceSubscriber{name: "router"}) + manager.Subscribe(&testNamespaceSubscriber{name: "router"}) + manager.Subscribe(&testNamespaceSubscriber{name: "buildermgr"}) + manager.Subscribe(&testNamespaceSubscriber{name: "router"}) expected := []string{"buildermgr", "router"} if !reflect.DeepEqual(expected, manager.SnapshotSubscribers()) { @@ -216,3 +222,46 @@ func TestNewBootstrappedNamespaceManager(t *testing.T) { t.Fatalf("expected env source, got %s", record.Source) } } + +func TestNamespaceManagerDispatchAdd(t *testing.T) { + manager := NewNamespaceManager() + manager.Upsert(NamespaceEvent{Type: NamespaceEventAdd, Name: "tenant-a", Source: NamespaceSourceWatcher}) + router := &testNamespaceSubscriber{name: "router"} + builder := &testNamespaceSubscriber{name: "buildermgr"} + manager.Subscribe(router) + manager.Subscribe(builder) + + record, ok, err := manager.DispatchAdd(context.Background(), "tenant-a") + if !ok || err != nil { + t.Fatalf("expected dispatch add success, ok=%v err=%v", ok, err) + } + if router.addCalls != 1 || builder.addCalls != 1 { + t.Fatalf("expected both subscribers to receive add call") + } + if record.Phase != NamespacePhaseActive { + t.Fatalf("expected active phase after successful dispatch, got %s", record.Phase) + } +} + +func TestNamespaceManagerDispatchResyncFailure(t *testing.T) { + manager := NewNamespaceManager() + manager.Upsert(NamespaceEvent{Type: NamespaceEventAdd, Name: "tenant-a", Source: NamespaceSourceWatcher}) + router := &testNamespaceSubscriber{name: "router"} + builder := &testNamespaceSubscriber{name: "buildermgr", resyncErr: errors.New("resync failed")} + manager.Subscribe(router) + manager.Subscribe(builder) + + record, ok, err := manager.DispatchResync(context.Background(), "tenant-a") + if !ok || err == nil { + t.Fatalf("expected dispatch resync failure, ok=%v err=%v", ok, err) + } + if builder.resyncCalls != 1 { + t.Fatalf("expected failing subscriber to receive resync call") + } + if record.Phase != NamespacePhaseFailed { + t.Fatalf("expected failed phase after dispatch error, got %s", record.Phase) + } + if record.RegisteredParts["buildermgr"].LastError != "resync failed" { + t.Fatalf("expected subscriber error to be stored") + } +}