Skip to content

Commit a50009f

Browse files
committed
feat(v12): replace Redis with Valkey, add alerts dedup/group cache, repin go-sdk everywhere
Three changes landing together because the last two depend on the first: - Redis -> Valkey: service, image (valkey/valkey:9-alpine), data dir and command in the installer; every consumer renamed to match (REDIS_ADDR/REDIS_PASSWORD/REDIS_DB -> VALKEY_ADDR/VALKEY_PASSWORD/ VALKEY_DB) in agent-manager and log-input. go-redis itself is untouched -- it already speaks Valkey's protocol. - plugins/alerts: wires go-sdk's new store/cache into isDuplicate and getPreviousAlertId as a fast path in front of ClickHouse, closing the race where a burst of alerts for the same GroupBy/DeduplicateBy value, arriving faster than ClickHouse's own write-visibility lag, each see "nothing yet" and each create their own root alert. An atomic Claim picks exactly one winner among concurrent callers instead. - go-sdk repinned to a596631 across every remaining plugin (ad-audit, aws, azure, bitdefender, crowdstrike, events, feeds, gcp, geolocation, o365, playground, rule-flood-guard, soar, soc-ai, sophos, stats) -- previously split across v1.1.27/28/31, none of which had this session's ClickHouse or cache fixes. go build/go vet/go test clean across all 17 plugin modules touched.
1 parent e34c500 commit a50009f

43 files changed

Lines changed: 262 additions & 223 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎agent-manager/authcache/authcache.go‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@ const (
2222
keyTTL = 30 * time.Minute
2323

2424
// Heals whatever the event-driven path missed: an eviction, a write that
25-
// failed, a Redis that was restarted.
25+
// failed, a Valkey that was restarted.
2626
republishEvery = 5 * time.Minute
2727

2828
opTimeout = 3 * time.Second
@@ -32,8 +32,8 @@ type Publisher struct {
3232
rdb *redis.Client
3333
}
3434

35-
// New returns nil when no Redis is configured, and every method on a nil
36-
// Publisher is a no-op — an install without Redis keeps working.
35+
// New returns nil when no Valkey is configured, and every method on a nil
36+
// Publisher is a no-op — an install without Valkey keeps working.
3737
func New(addr, password string, db int) *Publisher {
3838
if addr == "" {
3939
return nil

‎agent-manager/main.go‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -66,8 +66,8 @@ func main() {
6666
// Publishing keys where the log input reads them is what makes a deletion
6767
// take effect at once instead of when the input's own copy expires.
6868
agent.AuthCache = authcache.New(
69-
os.Getenv("REDIS_ADDR"),
70-
os.Getenv("REDIS_PASSWORD"),
69+
os.Getenv("VALKEY_ADDR"),
70+
os.Getenv("VALKEY_PASSWORD"),
7171
0,
7272
)
7373
defer func() { _ = agent.AuthCache.Close() }()

‎installer/docker/compose.go‎

Lines changed: 12 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -132,8 +132,8 @@ func (c *Compose) Populate(conf *config.Config, stack *StackConfig) error {
132132
"DB_PORT=5432",
133133
"DB_NAME=agentmanager",
134134
"PANEL_SERV_NAME=http://backend:8080",
135-
"REDIS_ADDR=redis:6379",
136-
"REDIS_PASSWORD=" + conf.Password,
135+
"VALKEY_ADDR=valkey:6379",
136+
"VALKEY_PASSWORD=" + conf.Password,
137137
},
138138
Logging: &dLogging,
139139
Deploy: &Deploy{
@@ -146,7 +146,7 @@ func (c *Compose) Populate(conf *config.Config, stack *StackConfig) error {
146146
},
147147
DependsOn: []string{
148148
"postgres",
149-
"redis",
149+
"valkey",
150150
},
151151
}
152152

@@ -398,24 +398,24 @@ func (c *Compose) Populate(conf *config.Config, stack *StackConfig) error {
398398
},
399399
}
400400

401-
redisMem := stack.ServiceResources["redis"].AssignedMemory
402-
c.Services["redis"] = Service{
403-
Image: utils.PointerOf[string]("redis:7-alpine"),
401+
valkeyMem := stack.ServiceResources["valkey"].AssignedMemory
402+
c.Services["valkey"] = Service{
403+
Image: utils.PointerOf[string]("valkey/valkey:9-alpine"),
404404
Volumes: []string{
405-
stack.RedisData + ":/data",
405+
stack.ValkeyData + ":/data",
406406
},
407407
// Holds the shared auth cache, which is rebuildable, so eviction under
408408
// pressure is preferable to refusing writes.
409409
Command: []string{
410-
"redis-server", "--requirepass", conf.Password,
410+
"valkey-server", "--requirepass", conf.Password,
411411
"--maxmemory-policy", "allkeys-lru",
412412
},
413413
Logging: &dLogging,
414414
Deploy: &Deploy{
415415
Placement: &pManager,
416416
Resources: &Resources{
417417
Limits: &Res{
418-
Memory: utils.PointerOf[string](fmt.Sprintf("%vM", redisMem)),
418+
Memory: utils.PointerOf[string](fmt.Sprintf("%vM", valkeyMem)),
419419
},
420420
},
421421
},
@@ -426,16 +426,16 @@ func (c *Compose) Populate(conf *config.Config, stack *StackConfig) error {
426426
Image: utils.PointerOf[string]("ghcr.io/utmstack/utmstack/log-input:${UTMSTACK_TAG}"),
427427
DependsOn: []string{
428428
"nats",
429-
"redis",
429+
"valkey",
430430
"agentmanager",
431431
},
432432
Volumes: []string{
433433
stack.Cert + ":/cert",
434434
},
435435
Environment: []string{
436436
"NATS_URL=nats://nats:4222",
437-
"REDIS_ADDR=redis:6379",
438-
"REDIS_PASSWORD=" + conf.Password,
437+
"VALKEY_ADDR=valkey:6379",
438+
"VALKEY_PASSWORD=" + conf.Password,
439439
"AGENT_MANAGER=agentmanager:9000",
440440
"BACKEND=http://backend:8080",
441441
"INTERNAL_KEY=" + conf.InternalKey,

‎installer/docker/plugins.go‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,12 @@ type PluginConfig struct {
2727
AgentManager string `yaml:"agentManager,omitempty"`
2828
Backend string `yaml:"backend,omitempty"`
2929
CertsFolder string `yaml:"certsFolder,omitempty"`
30+
Valkey ValkeyConfig `yaml:"valkey,omitempty"`
31+
}
32+
33+
type ValkeyConfig struct {
34+
Addr string `yaml:"addr,omitempty"`
35+
Password string `yaml:"password,omitempty"`
3036
}
3137

3238
type ClickHouseConfig struct {
@@ -95,6 +101,10 @@ func SetPluginsConfigs(conf *config.Config, stack *StackConfig) error {
95101
AgentManager: "agentmanager:9000",
96102
Backend: "http://backend:8080",
97103
CertsFolder: "/cert",
104+
Valkey: ValkeyConfig{
105+
Addr: "valkey:6379",
106+
Password: conf.Password,
107+
},
98108
}
99109

100110
clickHousePipeline := ClickHousePluginsConfig{}

‎installer/docker/stack.go‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ type StackConfig struct {
2323
ClickHouseConf string
2424
ClickHouseConfigD string
2525
NATSData string
26-
RedisData string
26+
ValkeyData string
2727
Cert string
2828
EventsEngineWorkdir string
2929
LocksDir string
@@ -67,7 +67,7 @@ func GetStackConfig() *StackConfig {
6767
// the cold-storage declaration, which is why both mount it.
6868
stackConfig.ClickHouseConfigD = utils.MakeDir(0777, cnf.DataDir, "clickhouse", "config.d")
6969
stackConfig.NATSData = utils.MakeDir(0777, cnf.DataDir, "nats")
70-
stackConfig.RedisData = utils.MakeDir(0777, cnf.DataDir, "redis")
70+
stackConfig.ValkeyData = utils.MakeDir(0777, cnf.DataDir, "valkey")
7171
stackConfig.LocksDir = utils.MakeDir(0777, cnf.DataDir, "locks")
7272
stackConfig.ShmFolder = utils.MakeDir(0777, cnf.DataDir, "tmpfs")
7373

@@ -76,7 +76,7 @@ func GetStackConfig() *StackConfig {
7676
{Name: "clickhouse", Priority: 1, MinMemory: 5120, MaxMemory: 60 * 1024},
7777
{Name: "log-input", Priority: 2, MinMemory: 256, MaxMemory: 1024},
7878
{Name: "nats", Priority: 3, MinMemory: 256, MaxMemory: 2 * 1024},
79-
{Name: "redis", Priority: 3, MinMemory: 128, MaxMemory: 1024},
79+
{Name: "valkey", Priority: 3, MinMemory: 128, MaxMemory: 1024},
8080
{Name: "backend", Priority: 3, MinMemory: 700, MaxMemory: 2 * 1024},
8181
{Name: "postgres", Priority: 2, MinMemory: 500, MaxMemory: 2 * 1024},
8282
{Name: "agentmanager", Priority: 3, MinMemory: 200, MaxMemory: 1024},

‎log-input/auth/auth.go‎

Lines changed: 10 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@ const (
2222
keyPrefixCollector = "auth:collector:"
2323
keyPrefixAPIKey = "auth:apikey:v2:"
2424

25-
redisTimeout = 3 * time.Second
25+
valkeyTimeout = 3 * time.Second
2626
)
2727

2828
type Service struct {
@@ -36,9 +36,9 @@ type Service struct {
3636
func New(cfg *config.Config) *Service {
3737
return &Service{
3838
rdb: redis.NewClient(&redis.Options{
39-
Addr: cfg.RedisAddr,
40-
Password: cfg.RedisPassword,
41-
DB: cfg.RedisDB,
39+
Addr: cfg.ValkeyAddr,
40+
Password: cfg.ValkeyPassword,
41+
DB: cfg.ValkeyDB,
4242
}),
4343
cfg: cfg,
4444
am: newAgentManagerClient(cfg),
@@ -49,7 +49,7 @@ func New(cfg *config.Config) *Service {
4949
func (s *Service) Close() error { return s.rdb.Close() }
5050

5151
func (s *Service) Ping(ctx context.Context) error {
52-
cCtx, cancel := context.WithTimeout(ctx, redisTimeout)
52+
cCtx, cancel := context.WithTimeout(ctx, valkeyTimeout)
5353
defer cancel()
5454
return s.rdb.Ping(cCtx).Err()
5555
}
@@ -129,10 +129,10 @@ func (s *Service) InternalKeyValid(key string) bool {
129129
return subtle.ConstantTimeCompare([]byte(key), []byte(s.cfg.InternalKey)) == 1
130130
}
131131

132-
// get reads a cached answer. An unreachable Redis reads as a miss so it cannot
133-
// authenticate anyone by accident.
132+
// get reads a cached answer. An unreachable Valkey reads as a miss so it
133+
// cannot authenticate anyone by accident.
134134
func (s *Service) get(ctx context.Context, key string) (string, bool) {
135-
cCtx, cancel := context.WithTimeout(ctx, redisTimeout)
135+
cCtx, cancel := context.WithTimeout(ctx, valkeyTimeout)
136136
defer cancel()
137137

138138
v, err := s.rdb.Get(cCtx, key).Result()
@@ -149,7 +149,7 @@ func (s *Service) get(ctx context.Context, key string) (string, bool) {
149149
}
150150

151151
func (s *Service) set(ctx context.Context, key, value string) {
152-
cCtx, cancel := context.WithTimeout(ctx, redisTimeout)
152+
cCtx, cancel := context.WithTimeout(ctx, valkeyTimeout)
153153
defer cancel()
154154

155155
if err := s.rdb.Set(cCtx, key, value, s.cfg.AuthTTL).Err(); err != nil {
@@ -169,7 +169,7 @@ func connectorPrefix(typ string) (string, bool) {
169169
return "", false
170170
}
171171

172-
// hashed keeps the API key out of Redis: what is cached is that it was
172+
// hashed keeps the API key out of Valkey: what is cached is that it was
173173
// accepted, not the key.
174174
func hashed(s string) string {
175175
sum := sha256.Sum256([]byte(s))

‎log-input/config/config.go‎

Lines changed: 25 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -28,9 +28,9 @@ type Config struct {
2828
// Lowering it strands messages already published to the higher subjects.
2929
Shards int
3030

31-
RedisAddr string
32-
RedisPassword string
33-
RedisDB int
31+
ValkeyAddr string
32+
ValkeyPassword string
33+
ValkeyDB int
3434

3535
AgentManager string
3636
Backend string
@@ -46,14 +46,14 @@ const (
4646
defaultListenAddr = "0.0.0.0:50051"
4747
defaultHealthAddr = "0.0.0.0:8080"
4848
defaultHTTPListenAddr = "0.0.0.0:50052"
49-
defaultShards = 16
50-
defaultTenant = "ce66672c-e36d-4761-a8c8-90058fee1a24"
51-
defaultAuthTTL = 5 * time.Minute
52-
defaultCertsFolder = "/cert"
53-
utmCertFileName = "utm.crt"
54-
utmCertFileKeyName = "utm.key"
55-
maxReasonableShards = 1024
56-
minAuthTTL = time.Second
49+
defaultShards = 16
50+
defaultTenant = "ce66672c-e36d-4761-a8c8-90058fee1a24"
51+
defaultAuthTTL = 5 * time.Minute
52+
defaultCertsFolder = "/cert"
53+
utmCertFileName = "utm.crt"
54+
utmCertFileKeyName = "utm.key"
55+
maxReasonableShards = 1024
56+
minAuthTTL = time.Second
5757
)
5858

5959
func Load() (*Config, error) {
@@ -63,27 +63,27 @@ func Load() (*Config, error) {
6363
ListenAddr: envOr("LISTEN_ADDR", defaultListenAddr),
6464
HealthAddr: envOr("HEALTH_ADDR", defaultHealthAddr),
6565
HTTPListenAddr: envOr("HTTP_LISTEN_ADDR", defaultHTTPListenAddr),
66-
CertFile: certs + "/" + utmCertFileName,
67-
KeyFile: certs + "/" + utmCertFileKeyName,
68-
NATSURL: os.Getenv("NATS_URL"),
69-
Shards: envInt("SHARDS", defaultShards),
70-
RedisAddr: os.Getenv("REDIS_ADDR"),
71-
RedisPassword: os.Getenv("REDIS_PASSWORD"),
72-
RedisDB: envInt("REDIS_DB", 0),
73-
AgentManager: os.Getenv("AGENT_MANAGER"),
74-
Backend: os.Getenv("BACKEND"),
75-
InternalKey: os.Getenv("INTERNAL_KEY"),
76-
DefaultTenant: envOr("DEFAULT_TENANT", defaultTenant),
77-
AuthTTL: envDuration("AUTH_TTL", defaultAuthTTL),
66+
CertFile: certs + "/" + utmCertFileName,
67+
KeyFile: certs + "/" + utmCertFileKeyName,
68+
NATSURL: os.Getenv("NATS_URL"),
69+
Shards: envInt("SHARDS", defaultShards),
70+
ValkeyAddr: os.Getenv("VALKEY_ADDR"),
71+
ValkeyPassword: os.Getenv("VALKEY_PASSWORD"),
72+
ValkeyDB: envInt("VALKEY_DB", 0),
73+
AgentManager: os.Getenv("AGENT_MANAGER"),
74+
Backend: os.Getenv("BACKEND"),
75+
InternalKey: os.Getenv("INTERNAL_KEY"),
76+
DefaultTenant: envOr("DEFAULT_TENANT", defaultTenant),
77+
AuthTTL: envDuration("AUTH_TTL", defaultAuthTTL),
7878
}
7979

8080
// Refused rather than defaulted: accepting a log with nowhere to put it
8181
// loses it.
8282
if c.NATSURL == "" {
8383
return nil, fmt.Errorf("NATS_URL is required")
8484
}
85-
if c.RedisAddr == "" {
86-
return nil, fmt.Errorf("REDIS_ADDR is required")
85+
if c.ValkeyAddr == "" {
86+
return nil, fmt.Errorf("VALKEY_ADDR is required")
8787
}
8888
if c.AgentManager == "" {
8989
return nil, fmt.Errorf("AGENT_MANAGER is required")

‎log-input/main.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,7 @@ func main() {
3434
if err := authService.Ping(ctx); err != nil {
3535
_ = catcher.Error("cannot reach the auth cache", err, map[string]any{
3636
"process": processName,
37-
"redis": cfg.RedisAddr,
37+
"valkey": cfg.ValkeyAddr,
3838
})
3939
time.Sleep(5 * time.Second)
4040
os.Exit(1)

‎plugins/ad-audit/go.mod‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ module github.com/utmstack/UTMStack/plugins/ad-audit
22

33
go 1.25.5
44

5-
require github.com/threatwinds/go-sdk v1.1.27-0.20260811073440-251cb9d842cd
5+
require github.com/threatwinds/go-sdk v1.1.37-0.20261001164002-a596631a84c5
66

77
require (
88
cel.dev/expr v0.25.2 // indirect

‎plugins/ad-audit/go.sum‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -88,8 +88,8 @@ github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXl
8888
github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY=
8989
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
9090
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
91-
github.com/threatwinds/go-sdk v1.1.27-0.20260811073440-251cb9d842cd h1:u5x6DfEkBErQurxPjtFYnd2xfW1ruD6kPxnYzGaHQCc=
92-
github.com/threatwinds/go-sdk v1.1.27-0.20260811073440-251cb9d842cd/go.mod h1:Mr5r+NTTe8k28gVS7EgEorn76VPtf0vPSaUCD2bWRSU=
91+
github.com/threatwinds/go-sdk v1.1.37-0.20261001164002-a596631a84c5 h1:AdoQX9Lo5Rb8e0FaPFIFZ607I7BfIJS9qBFKvWivJwU=
92+
github.com/threatwinds/go-sdk v1.1.37-0.20261001164002-a596631a84c5/go.mod h1:SnD2dR6FZOFZP+tnyTEdLT+EnQrrUH/pNaJIpU90xtk=
9393
github.com/tidwall/gjson v1.19.0 h1:xwxm7n691Uf3u5OFjzngavjGTh55KX5q/9w9xHW88JU=
9494
github.com/tidwall/gjson v1.19.0/go.mod h1:V37/opeE/JbLUOfH0QTXiNez2l0RUjYUhpT4szFQAfc=
9595
github.com/tidwall/match v1.2.0 h1:0pt8FlkOwjN2fPt4bIl4BoNxb98gGHN2ObFEDkrfZnM=

0 commit comments

Comments
 (0)