summaryrefslogtreecommitdiff
path: root/internal/controlplane/manager.go
diff options
context:
space:
mode:
Diffstat (limited to '')
-rw-r--r--internal/controlplane/manager.go38
1 files changed, 37 insertions, 1 deletions
diff --git a/internal/controlplane/manager.go b/internal/controlplane/manager.go
index 212963b..58c31fe 100644
--- a/internal/controlplane/manager.go
+++ b/internal/controlplane/manager.go
@@ -15,11 +15,14 @@ import (
const broadcastQueueSize = 128
+const redisHealthCheckTimeout = 500 * time.Millisecond
+
type managerStore interface {
LoadSnapshot(context.Context) (Snapshot, error)
DatabaseGeneration(context.Context) (int64, error)
PublishChange(context.Context, ChangeEvent) error
Subscribe(context.Context) (<-chan ChangeMessage, func() error, error)
+ PingRedis(context.Context) error
RedisEnabled() bool
}
@@ -129,7 +132,7 @@ func (m *Manager) RedisConnected() bool {
func (m *Manager) Run(ctx context.Context) {
var workers sync.WaitGroup
if m.store.RedisEnabled() {
- workers.Add(2)
+ workers.Add(3)
go func() {
defer workers.Done()
m.runSubscriptions(ctx)
@@ -138,6 +141,10 @@ func (m *Manager) Run(ctx context.Context) {
defer workers.Done()
m.runBroadcasts(ctx)
}()
+ go func() {
+ defer workers.Done()
+ m.runRedisHealth(ctx)
+ }()
} else {
m.logger.Info("control_plane_redis_disabled", "fallback", "postgres_polling")
}
@@ -145,6 +152,35 @@ func (m *Manager) Run(ctx context.Context) {
workers.Wait()
}
+func (m *Manager) runRedisHealth(ctx context.Context) {
+ interval := m.pollInterval
+ if interval > time.Second {
+ interval = time.Second
+ }
+ if interval <= 0 {
+ interval = time.Second
+ }
+ ticker := time.NewTicker(interval)
+ defer ticker.Stop()
+ for {
+ pingContext, cancel := context.WithTimeout(ctx, redisHealthCheckTimeout)
+ err := m.store.PingRedis(pingContext)
+ cancel()
+ if err != nil {
+ if m.redisConnected.Swap(false) {
+ m.logger.Warn("control_plane_redis_unavailable", "error", err, "fallback", "postgres_polling")
+ }
+ } else if !m.redisConnected.Swap(true) {
+ m.logger.Info("control_plane_redis_recovered")
+ }
+ select {
+ case <-ctx.Done():
+ return
+ case <-ticker.C:
+ }
+ }
+}
+
func (m *Manager) runPolling(ctx context.Context) {
ticker := time.NewTicker(m.pollInterval)
defer ticker.Stop()