mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-11 09:46:37 +08:00
fix(redis): add maintenance notification startup switch
Default Redis maintenance notification negotiation to disabled and apply the startup-only setting to both platform and Asynq clients.
This commit is contained in:
+14
-31
@@ -1,22 +1,10 @@
|
|||||||
# ──────────────────────────────────────────────────────────────────────────────
|
# ──────────────────────────────────────────────────────────────────────────────
|
||||||
# wavelet — 环境变量配置模板
|
# wavelet — 环境变量配置模板
|
||||||
# 复制此文件为 .env 并填入实际值: cp .env.example .env
|
# 复制此文件为 .env 并填入实际值: cp .env.example .env
|
||||||
# 环境变量优先级高于 config.yaml / config.docker.yaml
|
# 环境变量优先级高于 config.yaml
|
||||||
|
# docker compose 会读取本文件(env_file: .env)并替换 compose 中的 ${VAR}
|
||||||
# ──────────────────────────────────────────────────────────────────────────────
|
# ──────────────────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
# ─── Docker Compose 服务端口映射 ───────────────────────────────────────────────
|
|
||||||
APP_PORT=8000
|
|
||||||
POSTGRES_PORT=5432
|
|
||||||
REDIS_PORT=6379
|
|
||||||
JAEGER_UI_PORT=16686
|
|
||||||
JAEGER_OTLP_GRPC_PORT=4317
|
|
||||||
JAEGER_OTLP_HTTP_PORT=4318
|
|
||||||
|
|
||||||
# ─── PostgreSQL 容器配置(仅 docker-compose 使用)────────────────────────────
|
|
||||||
POSTGRES_DB=wavelet
|
|
||||||
POSTGRES_USER=postgres
|
|
||||||
POSTGRES_PASSWORD=postgres
|
|
||||||
|
|
||||||
# ─── 时区 ─────────────────────────────────────────────────────────────────────
|
# ─── 时区 ─────────────────────────────────────────────────────────────────────
|
||||||
TZ=Asia/Shanghai
|
TZ=Asia/Shanghai
|
||||||
|
|
||||||
@@ -35,7 +23,7 @@ APP_SESSION_HTTP_ONLY=true
|
|||||||
# HTTPS 部署时设为 true,HTTP 环境必须为 false
|
# HTTPS 部署时设为 true,HTTP 环境必须为 false
|
||||||
APP_SESSION_SECURE=true
|
APP_SESSION_SECURE=true
|
||||||
|
|
||||||
# ─── 数据库 ────────────────────────────────────────────────────────────────────
|
# ─── 数据库(PostgreSQL)──────────────────────────────────────────────────────
|
||||||
# 设置 DB_HOST 后自动启用 PostgreSQL,也可通过 DB_ENABLED 显式控制
|
# 设置 DB_HOST 后自动启用 PostgreSQL,也可通过 DB_ENABLED 显式控制
|
||||||
# DB_ENABLED=false 时使用 SQLite 作为后备数据库
|
# DB_ENABLED=false 时使用 SQLite 作为后备数据库
|
||||||
DB_ENABLED=true
|
DB_ENABLED=true
|
||||||
@@ -51,7 +39,7 @@ DB_TIMEZONE=Asia/Shanghai
|
|||||||
# DB_MAX_IDLE_CONN=16
|
# DB_MAX_IDLE_CONN=16
|
||||||
# DB_MAX_OPEN_CONN=128
|
# DB_MAX_OPEN_CONN=128
|
||||||
|
|
||||||
# ─── Redis ─────────────────────────────────────────────────────────────────────
|
# ─── Redis / Valkey ────────────────────────────────────────────────────────────
|
||||||
# 设置 REDIS_ADDR 后自动启用,也可通过 REDIS_ENABLED 显式控制
|
# 设置 REDIS_ADDR 后自动启用,也可通过 REDIS_ENABLED 显式控制
|
||||||
REDIS_ENABLED=true
|
REDIS_ENABLED=true
|
||||||
REDIS_ADDR=redis:6379
|
REDIS_ADDR=redis:6379
|
||||||
@@ -60,6 +48,10 @@ REDIS_ADDR=redis:6379
|
|||||||
# REDIS_DB=0
|
# REDIS_DB=0
|
||||||
REDIS_KEY_PREFIX=wavelet:
|
REDIS_KEY_PREFIX=wavelet:
|
||||||
# REDIS_POOL_SIZE=100
|
# REDIS_POOL_SIZE=100
|
||||||
|
# 启动时开关;修改后需重启服务
|
||||||
|
REDIS_MAINT_NOTIFICATIONS=false
|
||||||
|
# compose 宿主机映射端口(仅 docker-compose 使用)
|
||||||
|
# REDIS_PORT=6379
|
||||||
|
|
||||||
# ─── ClickHouse(可选,默认关闭)──────────────────────────────────────────
|
# ─── ClickHouse(可选,默认关闭)──────────────────────────────────────────
|
||||||
# 设置 CLICKHOUSE_HOST 后自动启用,也可显式控制
|
# 设置 CLICKHOUSE_HOST 后自动启用,也可显式控制
|
||||||
@@ -79,24 +71,15 @@ LOG_OUTPUT=stdout
|
|||||||
OTEL_EXPORTER_OTLP_ENDPOINT=http://jaeger:4317
|
OTEL_EXPORTER_OTLP_ENDPOINT=http://jaeger:4317
|
||||||
OTEL_EXPORTER_OTLP_INSECURE=true
|
OTEL_EXPORTER_OTLP_INSECURE=true
|
||||||
# 设为 0 关闭 tracing;本地 Jaeger 调试建议设为 1.0
|
# 设为 0 关闭 tracing;本地 Jaeger 调试建议设为 1.0
|
||||||
OTEL_SAMPLING_RATE=1.0
|
OTEL_SAMPLING_RATE=0.0
|
||||||
# 全局 Tracer 命名空间,默认为 github.com/Rain-kl/Wavelet
|
# 全局 Tracer 命名空间,默认为 github.com/Rain-kl/Wavelet
|
||||||
# OTEL_TRACER_NAME=github.com/Rain-kl/Wavelet
|
# OTEL_TRACER_NAME=github.com/Rain-kl/Wavelet
|
||||||
|
# compose 可选端口覆盖
|
||||||
# ─── S3 兼容存储(可选,默认关闭)──────────────────────────────────────────
|
# JAEGER_VERSION=2.19.0
|
||||||
# S3_ENABLED=false
|
# JAEGER_UI_PORT=16686
|
||||||
# S3_ENDPOINT=https://<account-id>.r2.cloudflarestorage.com
|
# JAEGER_OTLP_GRPC_PORT=4317
|
||||||
# S3_REGION=auto
|
# JAEGER_OTLP_HTTP_PORT=4318
|
||||||
# S3_BUCKET=<bucket-name>
|
|
||||||
# S3_ACCESS_KEY_ID=<access-key>
|
|
||||||
# S3_SECRET_ACCESS_KEY=<secret-key>
|
|
||||||
# S3_PATH_STYLE=false
|
|
||||||
# S3_CDN_URL=
|
|
||||||
|
|
||||||
# ─── Worker ────────────────────────────────────────────────────────────────────
|
# ─── Worker ────────────────────────────────────────────────────────────────────
|
||||||
# WORKER_CONCURRENCY=20
|
# WORKER_CONCURRENCY=20
|
||||||
# WORKER_STRICT_PRIORITY=false
|
# WORKER_STRICT_PRIORITY=false
|
||||||
|
|
||||||
# ─── Scheduler ─────────────────────────────────────────────────────────────────
|
|
||||||
# 未设置时默认为 @daily
|
|
||||||
# SCHEDULER_CLEANUP_CRON=0 */2 * * *
|
|
||||||
|
|||||||
@@ -67,6 +67,7 @@ redis:
|
|||||||
max_retries: 3
|
max_retries: 3
|
||||||
pool_timeout: 4
|
pool_timeout: 4
|
||||||
conn_max_idle_time: 300
|
conn_max_idle_time: 300
|
||||||
|
maint_notifications: false # 启动时开关;启用 Redis maintenance notifications 自动协商,修改后需重启
|
||||||
|
|
||||||
# ─── Logging ────────────────────────────────────────────────────────────────────
|
# ─── Logging ────────────────────────────────────────────────────────────────────
|
||||||
log:
|
log:
|
||||||
|
|||||||
@@ -208,6 +208,7 @@ func applyEnvOverrides(c *configModel) {
|
|||||||
c.Redis.DB = envInt("REDIS_DB", c.Redis.DB)
|
c.Redis.DB = envInt("REDIS_DB", c.Redis.DB)
|
||||||
c.Redis.KeyPrefix = envStr("REDIS_KEY_PREFIX", c.Redis.KeyPrefix)
|
c.Redis.KeyPrefix = envStr("REDIS_KEY_PREFIX", c.Redis.KeyPrefix)
|
||||||
c.Redis.PoolSize = envInt("REDIS_POOL_SIZE", c.Redis.PoolSize)
|
c.Redis.PoolSize = envInt("REDIS_POOL_SIZE", c.Redis.PoolSize)
|
||||||
|
c.Redis.MaintNotifications = envBool("REDIS_MAINT_NOTIFICATIONS", c.Redis.MaintNotifications)
|
||||||
|
|
||||||
// ─── ClickHouse ───
|
// ─── ClickHouse ───
|
||||||
if v, ok := os.LookupEnv("CLICKHOUSE_HOST"); ok {
|
if v, ok := os.LookupEnv("CLICKHOUSE_HOST"); ok {
|
||||||
|
|||||||
@@ -0,0 +1,14 @@
|
|||||||
|
package config
|
||||||
|
|
||||||
|
import "testing"
|
||||||
|
|
||||||
|
func TestApplyEnvOverridesRedisMaintNotifications(t *testing.T) {
|
||||||
|
t.Setenv("REDIS_MAINT_NOTIFICATIONS", "true")
|
||||||
|
|
||||||
|
cfg := &configModel{}
|
||||||
|
applyEnvOverrides(cfg)
|
||||||
|
|
||||||
|
if !cfg.Redis.MaintNotifications {
|
||||||
|
t.Fatal("REDIS_MAINT_NOTIFICATIONS=true was not applied")
|
||||||
|
}
|
||||||
|
}
|
||||||
+17
-16
@@ -87,22 +87,23 @@ type clickHouseConfig struct {
|
|||||||
|
|
||||||
// redisConfig Redis配置
|
// redisConfig Redis配置
|
||||||
type redisConfig struct {
|
type redisConfig struct {
|
||||||
Enabled bool `mapstructure:"enabled"`
|
Enabled bool `mapstructure:"enabled"`
|
||||||
Addrs []string `mapstructure:"addrs"`
|
Addrs []string `mapstructure:"addrs"`
|
||||||
Username string `mapstructure:"username"`
|
Username string `mapstructure:"username"`
|
||||||
Password string `mapstructure:"password"`
|
Password string `mapstructure:"password"`
|
||||||
DB int `mapstructure:"db"`
|
DB int `mapstructure:"db"`
|
||||||
ClusterMode bool `mapstructure:"cluster_mode"`
|
ClusterMode bool `mapstructure:"cluster_mode"`
|
||||||
MasterName string `mapstructure:"master_name"`
|
MasterName string `mapstructure:"master_name"`
|
||||||
KeyPrefix string `mapstructure:"key_prefix"`
|
KeyPrefix string `mapstructure:"key_prefix"`
|
||||||
PoolSize int `mapstructure:"pool_size"`
|
PoolSize int `mapstructure:"pool_size"`
|
||||||
MinIdleConn int `mapstructure:"min_idle_conn"`
|
MinIdleConn int `mapstructure:"min_idle_conn"`
|
||||||
DialTimeout int `mapstructure:"dial_timeout"`
|
DialTimeout int `mapstructure:"dial_timeout"`
|
||||||
ReadTimeout int `mapstructure:"read_timeout"`
|
ReadTimeout int `mapstructure:"read_timeout"`
|
||||||
WriteTimeout int `mapstructure:"write_timeout"`
|
WriteTimeout int `mapstructure:"write_timeout"`
|
||||||
MaxRetries int `mapstructure:"max_retries"`
|
MaxRetries int `mapstructure:"max_retries"`
|
||||||
PoolTimeout int `mapstructure:"pool_timeout"`
|
PoolTimeout int `mapstructure:"pool_timeout"`
|
||||||
ConnMaxIdleTime int `mapstructure:"conn_max_idle_time"`
|
ConnMaxIdleTime int `mapstructure:"conn_max_idle_time"`
|
||||||
|
MaintNotifications bool `mapstructure:"maint_notifications"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// logConfig 日志配置
|
// logConfig 日志配置
|
||||||
|
|||||||
+41
-32
@@ -35,45 +35,46 @@ func init() {
|
|||||||
if cfg.ClusterMode {
|
if cfg.ClusterMode {
|
||||||
// Cluster 模式
|
// Cluster 模式
|
||||||
Redis = redis.NewClusterClient(&redis.ClusterOptions{
|
Redis = redis.NewClusterClient(&redis.ClusterOptions{
|
||||||
Addrs: cfg.Addrs,
|
Addrs: cfg.Addrs,
|
||||||
Username: cfg.Username,
|
Username: cfg.Username,
|
||||||
Password: cfg.Password,
|
Password: cfg.Password,
|
||||||
PoolSize: cfg.PoolSize,
|
PoolSize: cfg.PoolSize,
|
||||||
MinIdleConns: cfg.MinIdleConn,
|
MinIdleConns: cfg.MinIdleConn,
|
||||||
DialTimeout: time.Duration(cfg.DialTimeout) * time.Second,
|
DialTimeout: time.Duration(cfg.DialTimeout) * time.Second,
|
||||||
ReadTimeout: time.Duration(cfg.ReadTimeout) * time.Second,
|
ReadTimeout: time.Duration(cfg.ReadTimeout) * time.Second,
|
||||||
WriteTimeout: time.Duration(cfg.WriteTimeout) * time.Second,
|
WriteTimeout: time.Duration(cfg.WriteTimeout) * time.Second,
|
||||||
MaxRetries: cfg.MaxRetries,
|
MaxRetries: cfg.MaxRetries,
|
||||||
PoolTimeout: time.Duration(cfg.PoolTimeout) * time.Second,
|
PoolTimeout: time.Duration(cfg.PoolTimeout) * time.Second,
|
||||||
ConnMaxIdleTime: time.Duration(cfg.ConnMaxIdleTime) * time.Second,
|
ConnMaxIdleTime: time.Duration(cfg.ConnMaxIdleTime) * time.Second,
|
||||||
MaintNotificationsConfig: &maintnotifications.Config{
|
MaintNotificationsConfig: redisMaintNotificationsConfig(cfg.MaintNotifications),
|
||||||
Mode: maintnotifications.ModeDisabled,
|
|
||||||
},
|
|
||||||
})
|
})
|
||||||
log.Println("[Redis] initialized in Cluster mode")
|
log.Println("[Redis] initialized in Cluster mode")
|
||||||
} else {
|
} else {
|
||||||
// Standalone 或 Sentinel 模式
|
// Standalone 或 Sentinel 模式
|
||||||
Redis = redis.NewUniversalClient(&redis.UniversalOptions{
|
options := &redis.UniversalOptions{
|
||||||
Addrs: cfg.Addrs,
|
Addrs: cfg.Addrs,
|
||||||
MasterName: cfg.MasterName, // 非空时启用 Sentinel
|
MasterName: cfg.MasterName, // 非空时启用 Sentinel
|
||||||
Username: cfg.Username,
|
Username: cfg.Username,
|
||||||
Password: cfg.Password,
|
Password: cfg.Password,
|
||||||
DB: cfg.DB,
|
DB: cfg.DB,
|
||||||
PoolSize: cfg.PoolSize,
|
PoolSize: cfg.PoolSize,
|
||||||
MinIdleConns: cfg.MinIdleConn,
|
MinIdleConns: cfg.MinIdleConn,
|
||||||
DialTimeout: time.Duration(cfg.DialTimeout) * time.Second,
|
DialTimeout: time.Duration(cfg.DialTimeout) * time.Second,
|
||||||
ReadTimeout: time.Duration(cfg.ReadTimeout) * time.Second,
|
ReadTimeout: time.Duration(cfg.ReadTimeout) * time.Second,
|
||||||
WriteTimeout: time.Duration(cfg.WriteTimeout) * time.Second,
|
WriteTimeout: time.Duration(cfg.WriteTimeout) * time.Second,
|
||||||
MaxRetries: cfg.MaxRetries,
|
MaxRetries: cfg.MaxRetries,
|
||||||
PoolTimeout: time.Duration(cfg.PoolTimeout) * time.Second,
|
PoolTimeout: time.Duration(cfg.PoolTimeout) * time.Second,
|
||||||
ConnMaxIdleTime: time.Duration(cfg.ConnMaxIdleTime) * time.Second,
|
ConnMaxIdleTime: time.Duration(cfg.ConnMaxIdleTime) * time.Second,
|
||||||
MaintNotificationsConfig: &maintnotifications.Config{
|
MaintNotificationsConfig: redisMaintNotificationsConfig(cfg.MaintNotifications),
|
||||||
Mode: maintnotifications.ModeDisabled,
|
}
|
||||||
},
|
|
||||||
})
|
|
||||||
if cfg.MasterName != "" {
|
if cfg.MasterName != "" {
|
||||||
|
client := redis.NewFailoverClient(options.Failover())
|
||||||
|
// FailoverOptions 暂不暴露该配置,在首次建连前写入客户端选项。
|
||||||
|
client.Options().MaintNotificationsConfig = redisMaintNotificationsConfig(cfg.MaintNotifications)
|
||||||
|
Redis = client
|
||||||
log.Println("[Redis] initialized in Sentinel mode")
|
log.Println("[Redis] initialized in Sentinel mode")
|
||||||
} else {
|
} else {
|
||||||
|
Redis = redis.NewUniversalClient(options)
|
||||||
log.Println("[Redis] initialized in Standalone mode")
|
log.Println("[Redis] initialized in Standalone mode")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -97,6 +98,14 @@ func init() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func redisMaintNotificationsConfig(enabled bool) *maintnotifications.Config {
|
||||||
|
mode := maintnotifications.ModeDisabled
|
||||||
|
if enabled {
|
||||||
|
mode = maintnotifications.ModeAuto
|
||||||
|
}
|
||||||
|
return &maintnotifications.Config{Mode: mode}
|
||||||
|
}
|
||||||
|
|
||||||
// PrefixedKey 返回带前缀的 Key
|
// PrefixedKey 返回带前缀的 Key
|
||||||
func PrefixedKey(key string) string {
|
func PrefixedKey(key string) string {
|
||||||
prefix := config.Config.Redis.KeyPrefix
|
prefix := config.Config.Redis.KeyPrefix
|
||||||
|
|||||||
@@ -0,0 +1,25 @@
|
|||||||
|
package db
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/redis/go-redis/v9/maintnotifications"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestRedisMaintNotificationsConfig(t *testing.T) {
|
||||||
|
for _, test := range []struct {
|
||||||
|
name string
|
||||||
|
enabled bool
|
||||||
|
want maintnotifications.Mode
|
||||||
|
}{
|
||||||
|
{name: "disabled by default", enabled: false, want: maintnotifications.ModeDisabled},
|
||||||
|
{name: "auto when enabled", enabled: true, want: maintnotifications.ModeAuto},
|
||||||
|
} {
|
||||||
|
t.Run(test.name, func(t *testing.T) {
|
||||||
|
cfg := redisMaintNotificationsConfig(test.enabled)
|
||||||
|
if cfg.Mode != test.want {
|
||||||
|
t.Fatalf("maintenance notifications mode = %v, want %v", cfg.Mode, test.want)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
+64
-16
@@ -7,8 +7,47 @@ package task
|
|||||||
import (
|
import (
|
||||||
"github.com/Rain-kl/Wavelet/internal/config"
|
"github.com/Rain-kl/Wavelet/internal/config"
|
||||||
"github.com/hibiken/asynq"
|
"github.com/hibiken/asynq"
|
||||||
|
"github.com/redis/go-redis/v9"
|
||||||
|
"github.com/redis/go-redis/v9/maintnotifications"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
type redisClientConnOpt struct {
|
||||||
|
options redis.Options
|
||||||
|
}
|
||||||
|
|
||||||
|
func (opt redisClientConnOpt) MakeRedisClient() interface{} {
|
||||||
|
return redis.NewClient(&opt.options)
|
||||||
|
}
|
||||||
|
|
||||||
|
type redisClusterConnOpt struct {
|
||||||
|
options redis.ClusterOptions
|
||||||
|
}
|
||||||
|
|
||||||
|
func (opt redisClusterConnOpt) MakeRedisClient() interface{} {
|
||||||
|
return redis.NewClusterClient(&opt.options)
|
||||||
|
}
|
||||||
|
|
||||||
|
type redisFailoverConnOpt struct {
|
||||||
|
options redis.FailoverOptions
|
||||||
|
maintNotificationsEnabled bool
|
||||||
|
}
|
||||||
|
|
||||||
|
func (opt redisFailoverConnOpt) MakeRedisClient() interface{} {
|
||||||
|
client := redis.NewFailoverClient(&opt.options)
|
||||||
|
// go-redis v9.16 does not expose maintenance notification settings on
|
||||||
|
// FailoverOptions, so apply the configured mode before the client is used.
|
||||||
|
client.Options().MaintNotificationsConfig = maintNotificationsConfig(opt.maintNotificationsEnabled)
|
||||||
|
return client
|
||||||
|
}
|
||||||
|
|
||||||
|
func maintNotificationsConfig(enabled bool) *maintnotifications.Config {
|
||||||
|
mode := maintnotifications.ModeDisabled
|
||||||
|
if enabled {
|
||||||
|
mode = maintnotifications.ModeAuto
|
||||||
|
}
|
||||||
|
return &maintnotifications.Config{Mode: mode}
|
||||||
|
}
|
||||||
|
|
||||||
// RedisOpt asynq Redis 连接配置(兼容 Standalone/Sentinel/Cluster)
|
// RedisOpt asynq Redis 连接配置(兼容 Standalone/Sentinel/Cluster)
|
||||||
var RedisOpt asynq.RedisConnOpt
|
var RedisOpt asynq.RedisConnOpt
|
||||||
|
|
||||||
@@ -26,20 +65,26 @@ func NewRedisConnOpt() asynq.RedisConnOpt {
|
|||||||
addrs := cfg.Addrs
|
addrs := cfg.Addrs
|
||||||
|
|
||||||
if cfg.ClusterMode {
|
if cfg.ClusterMode {
|
||||||
return asynq.RedisClusterClientOpt{
|
return redisClusterConnOpt{
|
||||||
Addrs: addrs,
|
options: redis.ClusterOptions{
|
||||||
Username: cfg.Username,
|
Addrs: addrs,
|
||||||
Password: cfg.Password,
|
Username: cfg.Username,
|
||||||
|
Password: cfg.Password,
|
||||||
|
MaintNotificationsConfig: maintNotificationsConfig(cfg.MaintNotifications),
|
||||||
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if cfg.MasterName != "" {
|
if cfg.MasterName != "" {
|
||||||
return asynq.RedisFailoverClientOpt{
|
return redisFailoverConnOpt{
|
||||||
MasterName: cfg.MasterName,
|
maintNotificationsEnabled: cfg.MaintNotifications,
|
||||||
SentinelAddrs: addrs,
|
options: redis.FailoverOptions{
|
||||||
Username: cfg.Username,
|
MasterName: cfg.MasterName,
|
||||||
Password: cfg.Password,
|
SentinelAddrs: addrs,
|
||||||
DB: cfg.DB,
|
Username: cfg.Username,
|
||||||
|
Password: cfg.Password,
|
||||||
|
DB: cfg.DB,
|
||||||
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -47,12 +92,15 @@ func NewRedisConnOpt() asynq.RedisConnOpt {
|
|||||||
if len(addrs) > 0 {
|
if len(addrs) > 0 {
|
||||||
addr = addrs[0]
|
addr = addrs[0]
|
||||||
}
|
}
|
||||||
return asynq.RedisClientOpt{
|
return redisClientConnOpt{
|
||||||
Addr: addr,
|
options: redis.Options{
|
||||||
Username: cfg.Username,
|
Addr: addr,
|
||||||
Password: cfg.Password,
|
Username: cfg.Username,
|
||||||
DB: cfg.DB,
|
Password: cfg.Password,
|
||||||
PoolSize: cfg.PoolSize,
|
DB: cfg.DB,
|
||||||
|
PoolSize: cfg.PoolSize,
|
||||||
|
MaintNotificationsConfig: maintNotificationsConfig(cfg.MaintNotifications),
|
||||||
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,73 @@
|
|||||||
|
package task
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/config"
|
||||||
|
"github.com/redis/go-redis/v9"
|
||||||
|
"github.com/redis/go-redis/v9/maintnotifications"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestNewRedisConnOptConfiguresMaintenanceNotifications(t *testing.T) {
|
||||||
|
previous := config.Config.Redis
|
||||||
|
t.Cleanup(func() { config.Config.Redis = previous })
|
||||||
|
|
||||||
|
for _, test := range []struct {
|
||||||
|
name string
|
||||||
|
enabled bool
|
||||||
|
want maintnotifications.Mode
|
||||||
|
}{
|
||||||
|
{name: "disabled by default", enabled: false, want: maintnotifications.ModeDisabled},
|
||||||
|
{name: "auto when enabled", enabled: true, want: maintnotifications.ModeAuto},
|
||||||
|
} {
|
||||||
|
t.Run(test.name, func(t *testing.T) {
|
||||||
|
config.Config.Redis.MaintNotifications = test.enabled
|
||||||
|
|
||||||
|
t.Run("standalone", func(t *testing.T) {
|
||||||
|
config.Config.Redis.ClusterMode = false
|
||||||
|
config.Config.Redis.MasterName = ""
|
||||||
|
config.Config.Redis.Addrs = []string{"127.0.0.1:6379"}
|
||||||
|
|
||||||
|
client, ok := NewRedisConnOpt().MakeRedisClient().(*redis.Client)
|
||||||
|
if !ok {
|
||||||
|
t.Fatal("standalone option did not create *redis.Client")
|
||||||
|
}
|
||||||
|
defer func() { _ = client.Close() }()
|
||||||
|
assertMaintenanceNotificationsMode(t, client.Options().MaintNotificationsConfig, test.want)
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("cluster", func(t *testing.T) {
|
||||||
|
config.Config.Redis.ClusterMode = true
|
||||||
|
config.Config.Redis.MasterName = ""
|
||||||
|
config.Config.Redis.Addrs = []string{"127.0.0.1:6379"}
|
||||||
|
|
||||||
|
client, ok := NewRedisConnOpt().MakeRedisClient().(*redis.ClusterClient)
|
||||||
|
if !ok {
|
||||||
|
t.Fatal("cluster option did not create *redis.ClusterClient")
|
||||||
|
}
|
||||||
|
defer func() { _ = client.Close() }()
|
||||||
|
assertMaintenanceNotificationsMode(t, client.Options().MaintNotificationsConfig, test.want)
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("sentinel", func(t *testing.T) {
|
||||||
|
config.Config.Redis.ClusterMode = false
|
||||||
|
config.Config.Redis.MasterName = "openflare"
|
||||||
|
config.Config.Redis.Addrs = []string{"127.0.0.1:26379"}
|
||||||
|
|
||||||
|
client, ok := NewRedisConnOpt().MakeRedisClient().(*redis.Client)
|
||||||
|
if !ok {
|
||||||
|
t.Fatal("sentinel option did not create *redis.Client")
|
||||||
|
}
|
||||||
|
defer func() { _ = client.Close() }()
|
||||||
|
assertMaintenanceNotificationsMode(t, client.Options().MaintNotificationsConfig, test.want)
|
||||||
|
})
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func assertMaintenanceNotificationsMode(t *testing.T, cfg *maintnotifications.Config, want maintnotifications.Mode) {
|
||||||
|
t.Helper()
|
||||||
|
if cfg == nil || cfg.Mode != want {
|
||||||
|
t.Fatalf("maintenance notifications mode = %v, want %v", cfg, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user