mirror of
https://github.com/openimsdk/open-im-server.git
synced 2026-09-15 06:41:43 +08:00
feat(redis): implement standalone gateway registration with Redis
This commit is contained in:
parent
ad2735a1cb
commit
ee672697c2
46
cmd/main.go
46
cmd/main.go
@ -37,6 +37,8 @@ import (
|
||||
"github.com/openimsdk/open-im-server/v3/pkg/authverify"
|
||||
"github.com/openimsdk/open-im-server/v3/pkg/common/config"
|
||||
"github.com/openimsdk/open-im-server/v3/pkg/common/prommetrics"
|
||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache/redis"
|
||||
"github.com/openimsdk/open-im-server/v3/pkg/dbbuild"
|
||||
"github.com/openimsdk/open-im-server/v3/version"
|
||||
"github.com/openimsdk/tools/discovery"
|
||||
"github.com/openimsdk/tools/discovery/inprocess"
|
||||
@ -458,16 +460,50 @@ func startRedisServerRegister(ctx context.Context, cfg *serverConfig, client dis
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
addr := net.JoinHostPort(registerIP, strconv.Itoa(apiPort))
|
||||
inprocess.SetLocalTarget(addr)
|
||||
timer := time.NewTimer(time.Second * 5)
|
||||
defer timer.Stop()
|
||||
const validTime = time.Second * 10
|
||||
dbb := dbbuild.NewBuilder(nil, &cfg.RedisConfig)
|
||||
rdb, err := dbb.Redis(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
gateway := redis.NewStandaloneGatewayRedis(rdb, validTime)
|
||||
selfAddr := net.JoinHostPort(registerIP, strconv.Itoa(apiPort))
|
||||
inprocess.SetLocalTarget(selfAddr)
|
||||
inprocess.SetBroadcastAddress(cfg.Share.Secret, func(ctx context.Context) ([]string, error) {
|
||||
address, err := gateway.GetGatewayAddrs(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
notSelf := make([]string, 0, len(address))
|
||||
for _, addr := range address {
|
||||
if addr != selfAddr {
|
||||
notSelf = append(notSelf, addr)
|
||||
}
|
||||
}
|
||||
return notSelf, nil
|
||||
})
|
||||
register := func() {
|
||||
ctx, cancel := context.WithTimeout(ctx, validTime/2)
|
||||
defer cancel()
|
||||
if err := gateway.RegisterGateway(ctx, selfAddr); err != nil {
|
||||
log.ZWarn(ctx, "gateway register failed", err, "address", selfAddr)
|
||||
}
|
||||
}
|
||||
timer := time.NewTimer(validTime / 2)
|
||||
defer func() {
|
||||
timer.Stop()
|
||||
ctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), time.Second)
|
||||
defer cancel()
|
||||
if err := gateway.UnregisterGateway(ctx, selfAddr); err != nil {
|
||||
log.ZWarn(ctx, "gateway unregister failed", err, "address", selfAddr)
|
||||
}
|
||||
}()
|
||||
for {
|
||||
select {
|
||||
case <-timer.C:
|
||||
register()
|
||||
case <-ctx.Done():
|
||||
return context.Cause(ctx)
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@ -9,6 +9,14 @@ imAdminUser:
|
||||
# Each entry here corresponds by index to the matching entry in the userIDs list above.
|
||||
nicknames: [superAdmin]
|
||||
|
||||
# queue: choose message queue engine
|
||||
# Supported values:
|
||||
# - kafka (default)
|
||||
# - redis
|
||||
# - memory (standalone only; microservices cannot use memory)
|
||||
queue: kafka
|
||||
|
||||
|
||||
# 1: For Android, iOS, Windows, Mac, and web platforms, only one instance can be online at a time
|
||||
multiLogin:
|
||||
policy: 1
|
||||
|
||||
54
pkg/common/storage/cache/redis/standalone_gateway.go
vendored
Normal file
54
pkg/common/storage/cache/redis/standalone_gateway.go
vendored
Normal file
@ -0,0 +1,54 @@
|
||||
package redis
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/redis/go-redis/v9"
|
||||
|
||||
"github.com/openimsdk/tools/errs"
|
||||
)
|
||||
|
||||
const standaloneGatewayHashKey = "STANDALONE_GATEWAY_REGISTRY"
|
||||
|
||||
type StandaloneGatewayRedis struct {
|
||||
rdb redis.UniversalClient
|
||||
validTime time.Duration
|
||||
}
|
||||
|
||||
func NewStandaloneGatewayRedis(rdb redis.UniversalClient, validTime time.Duration) *StandaloneGatewayRedis {
|
||||
return &StandaloneGatewayRedis{rdb: rdb, validTime: validTime}
|
||||
}
|
||||
|
||||
func (s *StandaloneGatewayRedis) RegisterGateway(ctx context.Context, addr string) error {
|
||||
pipe := s.rdb.Pipeline()
|
||||
pipe.HSet(ctx, standaloneGatewayHashKey, addr, strconv.FormatInt(time.Now().UnixMilli(), 10))
|
||||
pipe.Expire(ctx, standaloneGatewayHashKey, s.validTime*2)
|
||||
_, err := pipe.Exec(ctx)
|
||||
return errs.Wrap(err)
|
||||
}
|
||||
|
||||
func (s *StandaloneGatewayRedis) UnregisterGateway(ctx context.Context, addr string) error {
|
||||
return errs.Wrap(s.rdb.HDel(ctx, standaloneGatewayHashKey, addr).Err())
|
||||
}
|
||||
|
||||
func (s *StandaloneGatewayRedis) GetGatewayAddrs(ctx context.Context) ([]string, error) {
|
||||
gateways, err := s.rdb.HGetAll(ctx, standaloneGatewayHashKey).Result()
|
||||
if err != nil {
|
||||
return nil, errs.Wrap(err)
|
||||
}
|
||||
|
||||
now := time.Now()
|
||||
addrs := make([]string, 0, len(gateways))
|
||||
for addr, registeredAt := range gateways {
|
||||
registeredAtMs, err := strconv.ParseInt(registeredAt, 10, 64)
|
||||
if err != nil {
|
||||
return nil, errs.WrapMsg(err, "redis gateway register time is not int64", "addr", addr, "value", registeredAt)
|
||||
}
|
||||
if now.Sub(time.UnixMilli(registeredAtMs)) <= s.validTime {
|
||||
addrs = append(addrs, addr)
|
||||
}
|
||||
}
|
||||
return addrs, nil
|
||||
}
|
||||
52
pkg/common/storage/cache/redis/standalone_gateway_test.go
vendored
Normal file
52
pkg/common/storage/cache/redis/standalone_gateway_test.go
vendored
Normal file
@ -0,0 +1,52 @@
|
||||
package redis
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strconv"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/go-redis/redismock/v9"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestStandaloneGatewayRedisRegisterGateway(t *testing.T) {
|
||||
rdb, mock := redismock.NewClientMock()
|
||||
cache := NewStandaloneGatewayRedis(rdb, time.Second*10)
|
||||
|
||||
mock.Regexp().ExpectHSet(standaloneGatewayHashKey, "127.0.0.1:10001", `^[0-9]+$`).SetVal(1)
|
||||
mock.ExpectExpire(standaloneGatewayHashKey, time.Second*20).SetVal(true)
|
||||
|
||||
err := cache.RegisterGateway(context.Background(), "127.0.0.1:10001")
|
||||
require.NoError(t, err)
|
||||
assert.NoError(t, mock.ExpectationsWereMet())
|
||||
}
|
||||
|
||||
func TestStandaloneGatewayRedisUnregisterGateway(t *testing.T) {
|
||||
rdb, mock := redismock.NewClientMock()
|
||||
cache := NewStandaloneGatewayRedis(rdb, time.Second)
|
||||
|
||||
mock.ExpectHDel(standaloneGatewayHashKey, "127.0.0.1:10001").SetVal(1)
|
||||
|
||||
err := cache.UnregisterGateway(context.Background(), "127.0.0.1:10001")
|
||||
require.NoError(t, err)
|
||||
assert.NoError(t, mock.ExpectationsWereMet())
|
||||
}
|
||||
|
||||
func TestStandaloneGatewayRedisGetGatewayAddrs(t *testing.T) {
|
||||
rdb, mock := redismock.NewClientMock()
|
||||
cache := NewStandaloneGatewayRedis(rdb, time.Second*10)
|
||||
|
||||
now := time.Now()
|
||||
mock.ExpectHGetAll(standaloneGatewayHashKey).SetVal(map[string]string{
|
||||
"127.0.0.1:10001": strconv.FormatInt(now.Add(-time.Second).UnixMilli(), 10),
|
||||
"127.0.0.1:10002": strconv.FormatInt(now.Add(-time.Second*20).UnixMilli(), 10),
|
||||
"127.0.0.1:10003": strconv.FormatInt(now.Add(time.Second).UnixMilli(), 10),
|
||||
})
|
||||
|
||||
addrs, err := cache.GetGatewayAddrs(context.Background())
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, []string{"127.0.0.1:10001", "127.0.0.1:10003"}, addrs)
|
||||
assert.NoError(t, mock.ExpectationsWereMet())
|
||||
}
|
||||
Loading…
x
Reference in New Issue
Block a user