layer1: share watcher namespace manager bootstrap
This commit is contained in:
@@ -0,0 +1,19 @@
|
|||||||
|
# 2026-04-26 — NamespaceManager rewrite, step 37
|
||||||
|
|
||||||
|
## Цель шага
|
||||||
|
|
||||||
|
Убрать дублирование startup manager flow в трёх namespace watcher-ах.
|
||||||
|
|
||||||
|
## Что меняем
|
||||||
|
|
||||||
|
1. В `utils` добавляем helper `NewWatcherNamespaceManager()`.
|
||||||
|
2. Helper:
|
||||||
|
- создаёт `NamespaceManager`
|
||||||
|
- подписывает subscriber-ов
|
||||||
|
- выполняет `BootstrapAndDispatch()`
|
||||||
|
3. `buildermgr`, `router`, `executor/multitenant` используют новый helper.
|
||||||
|
|
||||||
|
## Что НЕ меняем
|
||||||
|
|
||||||
|
- не меняем semantics dispatch;
|
||||||
|
- не меняем runtime cleanup policy.
|
||||||
@@ -30,9 +30,8 @@ func StartNSWatcher(
|
|||||||
pkgw *packageWatcher,
|
pkgw *packageWatcher,
|
||||||
mgr manager.Interface,
|
mgr manager.Interface,
|
||||||
) {
|
) {
|
||||||
nsManager := utils.NewNamespaceManager()
|
nsManager, err := utils.NewWatcherNamespaceManager(ctx, utils.DefaultNSResolver().Snapshot(), utils.NamespaceSourceEnv, time.Now().UTC(), NewNamespaceSubscriber(envw, pkgw, mgr))
|
||||||
nsManager.Subscribe(NewNamespaceSubscriber(envw, pkgw, mgr))
|
if err != nil {
|
||||||
if _, err := nsManager.BootstrapAndDispatch(ctx, utils.DefaultNSResolver().Snapshot(), utils.NamespaceSourceEnv, time.Now().UTC()); err != nil {
|
|
||||||
logger.Error("buildermgr.NSWatcher: BootstrapAndDispatch failed", zap.Error(err))
|
logger.Error("buildermgr.NSWatcher: BootstrapAndDispatch failed", zap.Error(err))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -92,9 +92,8 @@ func StartNSWatcher(
|
|||||||
executorTypes map[fv1.ExecutorType]executortype.ExecutorType,
|
executorTypes map[fv1.ExecutorType]executortype.ExecutorType,
|
||||||
mgr manager.Interface,
|
mgr manager.Interface,
|
||||||
) {
|
) {
|
||||||
nsManager := utils.NewNamespaceManager()
|
nsManager, err := utils.NewWatcherNamespaceManager(ctx, utils.DefaultNSResolver().Snapshot(), utils.NamespaceSourceEnv, time.Now().UTC(), NewNamespaceSubscriber(logger, kubernetesClient, executorTypes, mgr))
|
||||||
nsManager.Subscribe(NewNamespaceSubscriber(logger, kubernetesClient, executorTypes, mgr))
|
if err != nil {
|
||||||
if _, err := nsManager.BootstrapAndDispatch(ctx, utils.DefaultNSResolver().Snapshot(), utils.NamespaceSourceEnv, time.Now().UTC()); err != nil {
|
|
||||||
logger.Error("multitenant.NSWatcher: BootstrapAndDispatch failed", zap.Error(err))
|
logger.Error("multitenant.NSWatcher: BootstrapAndDispatch failed", zap.Error(err))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -30,9 +30,8 @@ func StartNSWatcher(
|
|||||||
ts *HTTPTriggerSet,
|
ts *HTTPTriggerSet,
|
||||||
mgr manager.Interface,
|
mgr manager.Interface,
|
||||||
) {
|
) {
|
||||||
nsManager := utils.NewNamespaceManager()
|
nsManager, err := utils.NewWatcherNamespaceManager(ctx, utils.DefaultNSResolver().Snapshot(), utils.NamespaceSourceEnv, time.Now().UTC(), NewNamespaceSubscriber(ts, mgr))
|
||||||
nsManager.Subscribe(NewNamespaceSubscriber(ts, mgr))
|
if err != nil {
|
||||||
if _, err := nsManager.BootstrapAndDispatch(ctx, utils.DefaultNSResolver().Snapshot(), utils.NamespaceSourceEnv, time.Now().UTC()); err != nil {
|
|
||||||
logger.Error("router.NSWatcher: BootstrapAndDispatch failed", zap.Error(err))
|
logger.Error("router.NSWatcher: BootstrapAndDispatch failed", zap.Error(err))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -88,6 +88,18 @@ func NewBootstrappedNamespaceManager(resolver *NamespaceResolver, source Namespa
|
|||||||
return manager
|
return manager
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func NewWatcherNamespaceManager(ctx context.Context, namespaces []string, source NamespaceSource, observedAt time.Time, subscribers ...NamespaceSubscriber) (NamespaceManager, error) {
|
||||||
|
manager := NewNamespaceManager()
|
||||||
|
for _, subscriber := range subscribers {
|
||||||
|
if subscriber == nil {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
manager.Subscribe(subscriber)
|
||||||
|
}
|
||||||
|
_, err := manager.BootstrapAndDispatch(ctx, namespaces, source, observedAt)
|
||||||
|
return manager, err
|
||||||
|
}
|
||||||
|
|
||||||
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()
|
||||||
|
|||||||
@@ -270,6 +270,21 @@ func TestNewBootstrappedNamespaceManager(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestNewWatcherNamespaceManager(t *testing.T) {
|
||||||
|
router := &testNamespaceSubscriber{name: "router"}
|
||||||
|
builder := &testNamespaceSubscriber{name: "buildermgr"}
|
||||||
|
manager, err := NewWatcherNamespaceManager(context.Background(), []string{"tenant-b", "tenant-a"}, NamespaceSourceEnv, time.Now().UTC(), router, builder)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("expected watcher manager bootstrap success: %v", err)
|
||||||
|
}
|
||||||
|
if !reflect.DeepEqual([]string{"tenant-a", "tenant-b"}, manager.Snapshot()) {
|
||||||
|
t.Fatalf("expected watcher manager snapshot to match bootstrapped namespaces")
|
||||||
|
}
|
||||||
|
if router.addCalls != 2 || builder.addCalls != 2 {
|
||||||
|
t.Fatalf("expected all subscribers to receive bootstrap dispatch calls")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
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})
|
||||||
|
|||||||
Reference in New Issue
Block a user