summaryrefslogtreecommitdiff
path: root/internal/controlplane/manager_test.go
diff options
context:
space:
mode:
Diffstat (limited to 'internal/controlplane/manager_test.go')
-rw-r--r--internal/controlplane/manager_test.go49
1 files changed, 49 insertions, 0 deletions
diff --git a/internal/controlplane/manager_test.go b/internal/controlplane/manager_test.go
index b292060..78a0fdb 100644
--- a/internal/controlplane/manager_test.go
+++ b/internal/controlplane/manager_test.go
@@ -25,6 +25,7 @@ type fakeManagerStore struct {
publishErr error
published chan ChangeEvent
subscribe func(context.Context, int64) (<-chan ChangeMessage, func() error, error)
+ pingRedis func(context.Context) error
}
type capturePolicies struct{ values []domain.LimitPolicy }
@@ -84,6 +85,13 @@ func (s *fakeManagerStore) RedisEnabled() bool {
return s.redisEnabled
}
+func (s *fakeManagerStore) PingRedis(ctx context.Context) error {
+ if s.pingRedis != nil {
+ return s.pingRedis(ctx)
+ }
+ return nil
+}
+
func newTestManager(store managerStore, logger *slog.Logger, interval time.Duration) *Manager {
return NewManager(store, catalog.NewModels(nil), auth.NewDynamic(nil, false), logger, interval)
}
@@ -196,6 +204,14 @@ func TestSubscriptionMessageRestoresConnectedStateAfterPublishFailure(t *testing
store := newFakeManagerStore(1)
store.redisEnabled = true
store.publishErr = errors.New("redis unavailable")
+ var redisAvailable atomic.Bool
+ redisAvailable.Store(true)
+ store.pingRedis = func(context.Context) error {
+ if !redisAvailable.Load() {
+ return errors.New("redis unavailable")
+ }
+ return nil
+ }
messages := make(chan ChangeMessage, 1)
store.subscribe = func(_ context.Context, _ int64) (<-chan ChangeMessage, func() error, error) {
return messages, func() error { return nil }, nil
@@ -212,12 +228,14 @@ func TestSubscriptionMessageRestoresConnectedStateAfterPublishFailure(t *testing
close(done)
}()
waitUntil(t, time.Second, manager.RedisConnected)
+ redisAvailable.Store(false)
if err := manager.AfterMutation(context.Background(), 1, "model", "model-1"); err != nil {
t.Fatal(err)
}
waitUntil(t, time.Second, func() bool { return !manager.RedisConnected() })
store.snapshot.Store(Snapshot{Generation: 2})
+ redisAvailable.Store(true)
messages <- ChangeMessage{Payload: `{"generation":2,"resource":"model"}`}
waitUntil(t, time.Second, func() bool { return manager.RedisConnected() && manager.Generation() == 2 })
@@ -229,6 +247,37 @@ func TestSubscriptionMessageRestoresConnectedStateAfterPublishFailure(t *testing
}
}
+func TestRedisHealthCheckReportsFailureAndRecovery(t *testing.T) {
+ store := newFakeManagerStore(1)
+ store.redisEnabled = true
+ var redisAvailable atomic.Bool
+ redisAvailable.Store(true)
+ store.pingRedis = func(context.Context) error {
+ if !redisAvailable.Load() {
+ return errors.New("redis unavailable")
+ }
+ return nil
+ }
+ manager := newTestManager(store, slog.New(slog.NewTextHandler(&safeLogBuffer{}, nil)), 10*time.Millisecond)
+ ctx, cancel := context.WithCancel(context.Background())
+ done := make(chan struct{})
+ go func() {
+ manager.Run(ctx)
+ close(done)
+ }()
+ waitUntil(t, time.Second, manager.RedisConnected)
+ redisAvailable.Store(false)
+ waitUntil(t, time.Second, func() bool { return !manager.RedisConnected() })
+ redisAvailable.Store(true)
+ waitUntil(t, time.Second, manager.RedisConnected)
+ cancel()
+ select {
+ case <-done:
+ case <-time.After(time.Second):
+ t.Fatal("manager did not stop")
+ }
+}
+
func TestRedisCanBeDisabled(t *testing.T) {
store := newFakeManagerStore(3)
manager := newTestManager(store, slog.New(slog.NewTextHandler(&bytes.Buffer{}, nil)), 10*time.Millisecond)