layer1: add namespace remove dispatch step 33

This commit is contained in:
Naeel
2026-04-26 10:36:48 +03:00
parent 488157963a
commit d09bee3431
3 changed files with 83 additions and 1 deletions
+25
View File
@@ -57,6 +57,7 @@ type NamespaceManager interface {
SnapshotSubscribers() []string
Upsert(event NamespaceEvent) NamespaceRecord
DispatchAdd(ctx context.Context, namespace string) (NamespaceRecord, bool, error)
DispatchRemove(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)
@@ -250,6 +251,30 @@ func (m *inMemoryNamespaceManager) DispatchAdd(ctx context.Context, namespace st
})
}
func (m *inMemoryNamespaceManager) DispatchRemove(ctx context.Context, namespace string) (NamespaceRecord, bool, error) {
record, ok := m.Get(namespace)
if !ok {
return NamespaceRecord{}, false, nil
}
var firstErr error
for _, subscriber := range m.snapshotSubscriberObjects() {
err := subscriber.OnNamespaceRemove(ctx, record)
if err != nil && firstErr == nil {
firstErr = err
}
}
removedRecord := m.Upsert(NamespaceEvent{
Type: NamespaceEventRemove,
Name: record.Name,
Labels: record.Labels,
Source: record.Source,
ObservedAt: time.Now().UTC(),
})
return removedRecord, true, firstErr
}
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)
+41 -1
View File
@@ -11,8 +11,10 @@ import (
type testNamespaceSubscriber struct {
name string
addErr error
removeErr error
resyncErr error
addCalls int
removeCalls int
resyncCalls int
}
@@ -26,7 +28,8 @@ func (s *testNamespaceSubscriber) OnNamespaceAdd(ctx context.Context, record Nam
}
func (s *testNamespaceSubscriber) OnNamespaceRemove(ctx context.Context, record NamespaceRecord) error {
return nil
s.removeCalls++
return s.removeErr
}
func (s *testNamespaceSubscriber) OnNamespaceResync(ctx context.Context, record NamespaceRecord) error {
@@ -310,6 +313,43 @@ func TestNamespaceManagerDispatchResyncFailure(t *testing.T) {
}
}
func TestNamespaceManagerDispatchRemove(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.DispatchRemove(context.Background(), "tenant-a")
if !ok || err != nil {
t.Fatalf("expected dispatch remove success, ok=%v err=%v", ok, err)
}
if router.removeCalls != 1 || builder.removeCalls != 1 {
t.Fatalf("expected both subscribers to receive remove call")
}
if record.Phase != NamespacePhaseRemoved {
t.Fatalf("expected removed phase after dispatch remove, got %s", record.Phase)
}
}
func TestNamespaceManagerDispatchRemoveFailure(t *testing.T) {
manager := NewNamespaceManager()
manager.Upsert(NamespaceEvent{Type: NamespaceEventAdd, Name: "tenant-a", Source: NamespaceSourceWatcher})
router := &testNamespaceSubscriber{name: "router", removeErr: errors.New("remove failed")}
builder := &testNamespaceSubscriber{name: "buildermgr"}
manager.Subscribe(router)
manager.Subscribe(builder)
record, ok, err := manager.DispatchRemove(context.Background(), "tenant-a")
if !ok || err == nil {
t.Fatalf("expected dispatch remove failure, ok=%v err=%v", ok, err)
}
if record.Phase != NamespacePhaseRemoved {
t.Fatalf("expected removed phase even on dispatch remove error, got %s", record.Phase)
}
}
func TestNamespaceSubscriberFuncs(t *testing.T) {
addCalls := 0
removeCalls := 0