feat: Supports launching multiple monolithic instances, load-balanced by nginx, requiring only MongoDB and Redis at a minimum (#3759)
* fix: performance issues with Kafka caused by encapsulating the MQ interface * fix: admin token in standalone mode * fix: full id version * fix: resolve deadlock in cache eviction and improve GetBatch implementation * refactor: replace LongConn with ClientConn interface and simplify message handling * refactor: replace LongConn with ClientConn interface and simplify message handling * fix: seq use $setOnInsert for min_seq in conversation update * feat: add error code for handled friend requests and improve error handling in friend operations * refactor(msg): update regex pattern for conversationID to include a trailing colon * refactor: streamline cache initialization and enhance queue engine configuration * refactor: streamline cache initialization and enhance queue engine configuration * feat: add RegisterIP to API config and implement Redis server registration * feat(redis): implement standalone gateway registration with Redis * chore: update openimsdk/tools dependency to v0.0.50-alpha.121 * fix: change timer to ticker for improved periodic execution * refactor: replace Kafka topic configuration with mqbuild constants for improved readability * feat(redis): implement Redis-based locking mechanism for cron tasksdependabot/go_modules/golang.org/x/image-0.41.0
parent
7848d7c735
commit
10a577c310
@ -0,0 +1,14 @@
|
|||||||
|
package cron
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
lockLeaseTTL = time.Second * 300
|
||||||
|
)
|
||||||
|
|
||||||
|
type Locker interface {
|
||||||
|
ExecuteWithLock(ctx context.Context, taskName string, task func())
|
||||||
|
}
|
||||||
@ -0,0 +1,68 @@
|
|||||||
|
package cron
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/google/uuid"
|
||||||
|
"github.com/redis/go-redis/v9"
|
||||||
|
|
||||||
|
"github.com/openimsdk/tools/log"
|
||||||
|
)
|
||||||
|
|
||||||
|
func NewRedisLocker(client redis.UniversalClient) *RedisLocker {
|
||||||
|
return &RedisLocker{
|
||||||
|
client: client,
|
||||||
|
script: redis.NewScript(strings.TrimSpace(`
|
||||||
|
if redis.call("get", KEYS[1]) == ARGV[1] then
|
||||||
|
return redis.call("del", KEYS[1])
|
||||||
|
else
|
||||||
|
return 0
|
||||||
|
end
|
||||||
|
`)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type RedisLocker struct {
|
||||||
|
client redis.UniversalClient
|
||||||
|
script *redis.Script
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *RedisLocker) getKey(name string) string {
|
||||||
|
return "CRON_LOCKED:" + name
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *RedisLocker) lock(ctx context.Context, name string, owner string) (bool, error) {
|
||||||
|
ctx, cancel := context.WithTimeout(ctx, time.Second)
|
||||||
|
defer cancel()
|
||||||
|
return e.client.SetNX(ctx, e.getKey(name), owner, lockLeaseTTL).Result()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *RedisLocker) unlock(ctx context.Context, name string, owner string) error {
|
||||||
|
ctx, cancel := context.WithTimeout(ctx, time.Second)
|
||||||
|
defer cancel()
|
||||||
|
return e.script.Run(ctx, e.client, []string{e.getKey(name)}, owner).Err()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *RedisLocker) ExecuteWithLock(ctx context.Context, taskName string, task func()) {
|
||||||
|
owner := uuid.New().String()
|
||||||
|
ok, err := e.lock(ctx, taskName, owner)
|
||||||
|
if err != nil {
|
||||||
|
log.ZWarn(ctx, "cron lock get lock", err, "taskName", taskName)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
log.ZDebug(ctx, "cron lock get lock", "taskName", taskName, "ok", ok, "owner", owner)
|
||||||
|
if !ok {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
defer func() {
|
||||||
|
err := e.unlock(ctx, taskName, owner)
|
||||||
|
if err == nil {
|
||||||
|
log.ZDebug(ctx, "cron lock unlock", "taskName", taskName, "owner", owner)
|
||||||
|
} else {
|
||||||
|
log.ZWarn(ctx, "cron lock unlock", err, "taskName", taskName, "owner", owner)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
task()
|
||||||
|
}
|
||||||
@ -0,0 +1,41 @@
|
|||||||
|
package config
|
||||||
|
|
||||||
|
import (
|
||||||
|
"strings"
|
||||||
|
|
||||||
|
"github.com/openimsdk/tools/errs"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
QueueEngineKafka = "kafka"
|
||||||
|
QueueEngineRedis = "redis"
|
||||||
|
QueueEngineMemory = "memory"
|
||||||
|
)
|
||||||
|
|
||||||
|
func NormalizeQueueEngine(engine string) string {
|
||||||
|
switch strings.ToLower(strings.TrimSpace(engine)) {
|
||||||
|
case "", "kafka":
|
||||||
|
return QueueEngineKafka
|
||||||
|
case "redis":
|
||||||
|
return QueueEngineRedis
|
||||||
|
case "memory":
|
||||||
|
return QueueEngineMemory
|
||||||
|
default:
|
||||||
|
return strings.ToLower(strings.TrimSpace(engine))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func ValidateQueueEngine(engine string, standalone bool) (string, error) {
|
||||||
|
normalized := NormalizeQueueEngine(engine)
|
||||||
|
switch normalized {
|
||||||
|
case QueueEngineKafka, QueueEngineRedis:
|
||||||
|
return normalized, nil
|
||||||
|
case QueueEngineMemory:
|
||||||
|
if standalone {
|
||||||
|
return normalized, nil
|
||||||
|
}
|
||||||
|
return "", errs.ErrArgs.WrapMsg("unsupported queue engine for microservice deployment", "engine", engine)
|
||||||
|
default:
|
||||||
|
return "", errs.ErrArgs.WrapMsg("unsupported queue engine", "engine", engine)
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -0,0 +1,43 @@
|
|||||||
|
package config
|
||||||
|
|
||||||
|
import "encoding/json"
|
||||||
|
|
||||||
|
type EngineSelector struct {
|
||||||
|
Engine string `yaml:"engine" mapstructure:"engine" json:"engine"`
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e EngineSelector) String() string {
|
||||||
|
return e.Engine
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *EngineSelector) UnmarshalYAML(unmarshal func(any) error) error {
|
||||||
|
var engine string
|
||||||
|
if err := unmarshal(&engine); err == nil {
|
||||||
|
e.Engine = engine
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
var cfg struct {
|
||||||
|
Engine string `yaml:"engine"`
|
||||||
|
}
|
||||||
|
if err := unmarshal(&cfg); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
e.Engine = cfg.Engine
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *EngineSelector) UnmarshalJSON(data []byte) error {
|
||||||
|
var engine string
|
||||||
|
if err := json.Unmarshal(data, &engine); err == nil {
|
||||||
|
e.Engine = engine
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
var cfg struct {
|
||||||
|
Engine string `json:"engine"`
|
||||||
|
}
|
||||||
|
if err := json.Unmarshal(data, &cfg); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
e.Engine = cfg.Engine
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@ -1,50 +0,0 @@
|
|||||||
package mcache
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache/cachekey"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/database"
|
|
||||||
"github.com/openimsdk/tools/s3/minio"
|
|
||||||
)
|
|
||||||
|
|
||||||
func NewMinioCache(cache database.Cache) minio.Cache {
|
|
||||||
return &minioCache{
|
|
||||||
cache: cache,
|
|
||||||
expireTime: time.Hour * 24 * 7,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
type minioCache struct {
|
|
||||||
cache database.Cache
|
|
||||||
expireTime time.Duration
|
|
||||||
}
|
|
||||||
|
|
||||||
func (g *minioCache) getObjectImageInfoKey(key string) string {
|
|
||||||
return cachekey.GetObjectImageInfoKey(key)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (g *minioCache) getMinioImageThumbnailKey(key string, format string, width int, height int) string {
|
|
||||||
return cachekey.GetMinioImageThumbnailKey(key, format, width, height)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (g *minioCache) DelObjectImageInfoKey(ctx context.Context, keys ...string) error {
|
|
||||||
ks := make([]string, 0, len(keys))
|
|
||||||
for _, key := range keys {
|
|
||||||
ks = append(ks, g.getObjectImageInfoKey(key))
|
|
||||||
}
|
|
||||||
return g.cache.Del(ctx, ks)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (g *minioCache) DelImageThumbnailKey(ctx context.Context, key string, format string, width int, height int) error {
|
|
||||||
return g.cache.Del(ctx, []string{g.getMinioImageThumbnailKey(key, format, width, height)})
|
|
||||||
}
|
|
||||||
|
|
||||||
func (g *minioCache) GetImageObjectKeyInfo(ctx context.Context, key string, fn func(ctx context.Context) (*minio.ImageInfo, error)) (*minio.ImageInfo, error) {
|
|
||||||
return getCache[*minio.ImageInfo](ctx, g.cache, g.getObjectImageInfoKey(key), g.expireTime, fn)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (g *minioCache) GetThumbnailKey(ctx context.Context, key string, format string, width int, height int, minioCache func(ctx context.Context) (string, error)) (string, error) {
|
|
||||||
return getCache[string](ctx, g.cache, g.getMinioImageThumbnailKey(key, format, width, height), g.expireTime, minioCache)
|
|
||||||
}
|
|
||||||
@ -1,132 +0,0 @@
|
|||||||
package mcache
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"strconv"
|
|
||||||
"sync"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache/cachekey"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/database"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/model"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/localcache"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/localcache/lru"
|
|
||||||
"github.com/openimsdk/tools/errs"
|
|
||||||
"github.com/openimsdk/tools/utils/datautil"
|
|
||||||
"github.com/redis/go-redis/v9"
|
|
||||||
)
|
|
||||||
|
|
||||||
var (
|
|
||||||
memMsgCache lru.LRU[string, *model.MsgInfoModel]
|
|
||||||
initMemMsgCache sync.Once
|
|
||||||
)
|
|
||||||
|
|
||||||
func NewMsgCache(cache database.Cache, msgDocDatabase database.Msg) cache.MsgCache {
|
|
||||||
initMemMsgCache.Do(func() {
|
|
||||||
memMsgCache = lru.NewLazyLRU[string, *model.MsgInfoModel](1024*8, time.Hour, time.Second*10, localcache.EmptyTarget{}, nil)
|
|
||||||
})
|
|
||||||
return &msgCache{
|
|
||||||
cache: cache,
|
|
||||||
msgDocDatabase: msgDocDatabase,
|
|
||||||
memMsgCache: memMsgCache,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
type msgCache struct {
|
|
||||||
cache database.Cache
|
|
||||||
msgDocDatabase database.Msg
|
|
||||||
memMsgCache lru.LRU[string, *model.MsgInfoModel]
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *msgCache) getSendMsgKey(id string) string {
|
|
||||||
return cachekey.GetSendMsgKey(id)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *msgCache) SetSendMsgStatus(ctx context.Context, id string, status int32) error {
|
|
||||||
return x.cache.Set(ctx, x.getSendMsgKey(id), strconv.Itoa(int(status)), time.Hour*24)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *msgCache) GetSendMsgStatus(ctx context.Context, id string) (int32, error) {
|
|
||||||
key := x.getSendMsgKey(id)
|
|
||||||
res, err := x.cache.Get(ctx, []string{key})
|
|
||||||
if err != nil {
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
val, ok := res[key]
|
|
||||||
if !ok {
|
|
||||||
return 0, errs.Wrap(redis.Nil)
|
|
||||||
}
|
|
||||||
status, err := strconv.Atoi(val)
|
|
||||||
if err != nil {
|
|
||||||
return 0, errs.WrapMsg(err, "GetSendMsgStatus strconv.Atoi error", "val", val)
|
|
||||||
}
|
|
||||||
return int32(status), nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *msgCache) getMsgCacheKey(conversationID string, seq int64) string {
|
|
||||||
return cachekey.GetMsgCacheKey(conversationID, seq)
|
|
||||||
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *msgCache) GetMessageBySeqs(ctx context.Context, conversationID string, seqs []int64) ([]*model.MsgInfoModel, error) {
|
|
||||||
if len(seqs) == 0 {
|
|
||||||
return nil, nil
|
|
||||||
}
|
|
||||||
keys := make([]string, 0, len(seqs))
|
|
||||||
keySeq := make(map[string]int64, len(seqs))
|
|
||||||
for _, seq := range seqs {
|
|
||||||
key := x.getMsgCacheKey(conversationID, seq)
|
|
||||||
keys = append(keys, key)
|
|
||||||
keySeq[key] = seq
|
|
||||||
}
|
|
||||||
res, err := x.memMsgCache.GetBatch(keys, func(keys []string) (map[string]*model.MsgInfoModel, error) {
|
|
||||||
findSeqs := make([]int64, 0, len(keys))
|
|
||||||
for _, key := range keys {
|
|
||||||
seq, ok := keySeq[key]
|
|
||||||
if !ok {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
findSeqs = append(findSeqs, seq)
|
|
||||||
}
|
|
||||||
res, err := x.msgDocDatabase.FindSeqs(ctx, conversationID, seqs)
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
kv := make(map[string]*model.MsgInfoModel)
|
|
||||||
for i := range res {
|
|
||||||
msg := res[i]
|
|
||||||
if msg == nil || msg.Msg == nil || msg.Msg.Seq <= 0 {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
key := x.getMsgCacheKey(conversationID, msg.Msg.Seq)
|
|
||||||
kv[key] = msg
|
|
||||||
}
|
|
||||||
return kv, nil
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
return datautil.Values(res), nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x msgCache) DelMessageBySeqs(ctx context.Context, conversationID string, seqs []int64) error {
|
|
||||||
if len(seqs) == 0 {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
for _, seq := range seqs {
|
|
||||||
x.memMsgCache.Del(x.getMsgCacheKey(conversationID, seq))
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *msgCache) SetMessageBySeqs(ctx context.Context, conversationID string, msgs []*model.MsgInfoModel) error {
|
|
||||||
for i := range msgs {
|
|
||||||
msg := msgs[i]
|
|
||||||
if msg == nil || msg.Msg == nil || msg.Msg.Seq <= 0 {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
x.memMsgCache.Set(x.getMsgCacheKey(conversationID, msg.Msg.Seq), msg)
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
@ -1,82 +0,0 @@
|
|||||||
package mcache
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"sync"
|
|
||||||
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache"
|
|
||||||
)
|
|
||||||
|
|
||||||
var (
|
|
||||||
globalOnlineCache cache.OnlineCache
|
|
||||||
globalOnlineOnce sync.Once
|
|
||||||
)
|
|
||||||
|
|
||||||
func NewOnlineCache() cache.OnlineCache {
|
|
||||||
globalOnlineOnce.Do(func() {
|
|
||||||
globalOnlineCache = &onlineCache{
|
|
||||||
user: make(map[string]map[int32]struct{}),
|
|
||||||
}
|
|
||||||
})
|
|
||||||
return globalOnlineCache
|
|
||||||
}
|
|
||||||
|
|
||||||
type onlineCache struct {
|
|
||||||
lock sync.RWMutex
|
|
||||||
user map[string]map[int32]struct{}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *onlineCache) GetOnline(ctx context.Context, userID string) ([]int32, error) {
|
|
||||||
x.lock.RLock()
|
|
||||||
defer x.lock.RUnlock()
|
|
||||||
pSet, ok := x.user[userID]
|
|
||||||
if !ok {
|
|
||||||
return nil, nil
|
|
||||||
}
|
|
||||||
res := make([]int32, 0, len(pSet))
|
|
||||||
for k := range pSet {
|
|
||||||
res = append(res, k)
|
|
||||||
}
|
|
||||||
return res, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *onlineCache) SetUserOnline(ctx context.Context, userID string, online, offline []int32) error {
|
|
||||||
x.lock.Lock()
|
|
||||||
defer x.lock.Unlock()
|
|
||||||
pSet, ok := x.user[userID]
|
|
||||||
if ok {
|
|
||||||
for _, p := range offline {
|
|
||||||
delete(pSet, p)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if len(online) > 0 {
|
|
||||||
if !ok {
|
|
||||||
pSet = make(map[int32]struct{})
|
|
||||||
x.user[userID] = pSet
|
|
||||||
}
|
|
||||||
for _, p := range online {
|
|
||||||
pSet[p] = struct{}{}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if len(pSet) == 0 {
|
|
||||||
delete(x.user, userID)
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *onlineCache) GetAllOnlineUsers(ctx context.Context, cursor uint64) (map[string][]int32, uint64, error) {
|
|
||||||
if cursor != 0 {
|
|
||||||
return nil, 0, nil
|
|
||||||
}
|
|
||||||
x.lock.RLock()
|
|
||||||
defer x.lock.RUnlock()
|
|
||||||
res := make(map[string][]int32)
|
|
||||||
for k, v := range x.user {
|
|
||||||
pSet := make([]int32, 0, len(v))
|
|
||||||
for p := range v {
|
|
||||||
pSet = append(pSet, p)
|
|
||||||
}
|
|
||||||
res[k] = pSet
|
|
||||||
}
|
|
||||||
return res, 0, nil
|
|
||||||
}
|
|
||||||
@ -1,79 +0,0 @@
|
|||||||
package mcache
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/database"
|
|
||||||
)
|
|
||||||
|
|
||||||
func NewSeqConversationCache(sc database.SeqConversation) cache.SeqConversationCache {
|
|
||||||
return &seqConversationCache{
|
|
||||||
sc: sc,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
type seqConversationCache struct {
|
|
||||||
sc database.SeqConversation
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *seqConversationCache) Malloc(ctx context.Context, conversationID string, size int64) (int64, error) {
|
|
||||||
return x.sc.Malloc(ctx, conversationID, size)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *seqConversationCache) SetMinSeq(ctx context.Context, conversationID string, seq int64) error {
|
|
||||||
return x.sc.SetMinSeq(ctx, conversationID, seq)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *seqConversationCache) GetMinSeq(ctx context.Context, conversationID string) (int64, error) {
|
|
||||||
return x.sc.GetMinSeq(ctx, conversationID)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *seqConversationCache) GetMaxSeqs(ctx context.Context, conversationIDs []string) (map[string]int64, error) {
|
|
||||||
res := make(map[string]int64)
|
|
||||||
for _, conversationID := range conversationIDs {
|
|
||||||
seq, err := x.GetMaxSeq(ctx, conversationID)
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
res[conversationID] = seq
|
|
||||||
}
|
|
||||||
return res, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *seqConversationCache) GetMaxSeqsWithTime(ctx context.Context, conversationIDs []string) (map[string]database.SeqTime, error) {
|
|
||||||
res := make(map[string]database.SeqTime)
|
|
||||||
for _, conversationID := range conversationIDs {
|
|
||||||
seq, err := x.GetMaxSeq(ctx, conversationID)
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
res[conversationID] = database.SeqTime{Seq: seq}
|
|
||||||
}
|
|
||||||
return res, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *seqConversationCache) GetMaxSeq(ctx context.Context, conversationID string) (int64, error) {
|
|
||||||
return x.sc.GetMaxSeq(ctx, conversationID)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *seqConversationCache) GetMaxSeqWithTime(ctx context.Context, conversationID string) (database.SeqTime, error) {
|
|
||||||
seq, err := x.GetMinSeq(ctx, conversationID)
|
|
||||||
if err != nil {
|
|
||||||
return database.SeqTime{}, err
|
|
||||||
}
|
|
||||||
return database.SeqTime{Seq: seq}, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *seqConversationCache) SetMinSeqs(ctx context.Context, seqs map[string]int64) error {
|
|
||||||
for conversationID, seq := range seqs {
|
|
||||||
if err := x.sc.SetMinSeq(ctx, conversationID, seq); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *seqConversationCache) GetCacheMaxSeqWithTime(ctx context.Context, conversationIDs []string) (map[string]database.SeqTime, error) {
|
|
||||||
return x.GetMaxSeqsWithTime(ctx, conversationIDs)
|
|
||||||
}
|
|
||||||
@ -1,98 +0,0 @@
|
|||||||
package mcache
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"strconv"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache/cachekey"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/database"
|
|
||||||
"github.com/openimsdk/tools/errs"
|
|
||||||
"github.com/redis/go-redis/v9"
|
|
||||||
)
|
|
||||||
|
|
||||||
func NewThirdCache(cache database.Cache) cache.ThirdCache {
|
|
||||||
return &thirdCache{
|
|
||||||
cache: cache,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
type thirdCache struct {
|
|
||||||
cache database.Cache
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) getGetuiTokenKey() string {
|
|
||||||
return cachekey.GetGetuiTokenKey()
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) getGetuiTaskIDKey() string {
|
|
||||||
return cachekey.GetGetuiTaskIDKey()
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) getUserBadgeUnreadCountSumKey(userID string) string {
|
|
||||||
return cachekey.GetUserBadgeUnreadCountSumKey(userID)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) getFcmAccountTokenKey(account string, platformID int) string {
|
|
||||||
return cachekey.GetFcmAccountTokenKey(account, platformID)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) get(ctx context.Context, key string) (string, error) {
|
|
||||||
res, err := c.cache.Get(ctx, []string{key})
|
|
||||||
if err != nil {
|
|
||||||
return "", err
|
|
||||||
}
|
|
||||||
if val, ok := res[key]; ok {
|
|
||||||
return val, nil
|
|
||||||
}
|
|
||||||
return "", errs.Wrap(redis.Nil)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) SetFcmToken(ctx context.Context, account string, platformID int, fcmToken string, expireTime int64) (err error) {
|
|
||||||
return errs.Wrap(c.cache.Set(ctx, c.getFcmAccountTokenKey(account, platformID), fcmToken, time.Duration(expireTime)*time.Second))
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) GetFcmToken(ctx context.Context, account string, platformID int) (string, error) {
|
|
||||||
return c.get(ctx, c.getFcmAccountTokenKey(account, platformID))
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) DelFcmToken(ctx context.Context, account string, platformID int) error {
|
|
||||||
return c.cache.Del(ctx, []string{c.getFcmAccountTokenKey(account, platformID)})
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) IncrUserBadgeUnreadCountSum(ctx context.Context, userID string) (int, error) {
|
|
||||||
return c.cache.Incr(ctx, c.getUserBadgeUnreadCountSumKey(userID), 1)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) SetUserBadgeUnreadCountSum(ctx context.Context, userID string, value int) error {
|
|
||||||
return c.cache.Set(ctx, c.getUserBadgeUnreadCountSumKey(userID), strconv.Itoa(value), 0)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) GetUserBadgeUnreadCountSum(ctx context.Context, userID string) (int, error) {
|
|
||||||
str, err := c.get(ctx, c.getUserBadgeUnreadCountSumKey(userID))
|
|
||||||
if err != nil {
|
|
||||||
return 0, err
|
|
||||||
}
|
|
||||||
val, err := strconv.Atoi(str)
|
|
||||||
if err != nil {
|
|
||||||
return 0, errs.WrapMsg(err, "strconv.Atoi", "str", str)
|
|
||||||
}
|
|
||||||
return val, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) SetGetuiToken(ctx context.Context, token string, expireTime int64) error {
|
|
||||||
return c.cache.Set(ctx, c.getGetuiTokenKey(), token, time.Duration(expireTime)*time.Second)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) GetGetuiToken(ctx context.Context) (string, error) {
|
|
||||||
return c.get(ctx, c.getGetuiTokenKey())
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) SetGetuiTaskID(ctx context.Context, taskID string, expireTime int64) error {
|
|
||||||
return c.cache.Set(ctx, c.getGetuiTaskIDKey(), taskID, time.Duration(expireTime)*time.Second)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *thirdCache) GetGetuiTaskID(ctx context.Context) (string, error) {
|
|
||||||
return c.get(ctx, c.getGetuiTaskIDKey())
|
|
||||||
}
|
|
||||||
@ -1,166 +0,0 @@
|
|||||||
package mcache
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"fmt"
|
|
||||||
"strconv"
|
|
||||||
"strings"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/cache/cachekey"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/database"
|
|
||||||
"github.com/openimsdk/tools/errs"
|
|
||||||
"github.com/openimsdk/tools/log"
|
|
||||||
)
|
|
||||||
|
|
||||||
func NewTokenCacheModel(cache database.Cache, accessExpire int64) cache.TokenModel {
|
|
||||||
c := &tokenCache{cache: cache}
|
|
||||||
c.accessExpire = c.getExpireTime(accessExpire)
|
|
||||||
return c
|
|
||||||
}
|
|
||||||
|
|
||||||
type tokenCache struct {
|
|
||||||
cache database.Cache
|
|
||||||
accessExpire time.Duration
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *tokenCache) getTokenKey(userID string, platformID int, token string) string {
|
|
||||||
return cachekey.GetTokenKey(userID, platformID) + ":" + token
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *tokenCache) SetTokenFlag(ctx context.Context, userID string, platformID int, token string, flag int) error {
|
|
||||||
return x.cache.Set(ctx, x.getTokenKey(userID, platformID, token), strconv.Itoa(flag), x.accessExpire)
|
|
||||||
}
|
|
||||||
|
|
||||||
// SetTokenFlagEx set token and flag with expire time
|
|
||||||
func (x *tokenCache) SetTokenFlagEx(ctx context.Context, userID string, platformID int, token string, flag int) error {
|
|
||||||
return x.SetTokenFlag(ctx, userID, platformID, token, flag)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *tokenCache) GetTokensWithoutError(ctx context.Context, userID string, platformID int) (map[string]int, error) {
|
|
||||||
prefix := x.getTokenKey(userID, platformID, "")
|
|
||||||
m, err := x.cache.Prefix(ctx, prefix)
|
|
||||||
if err != nil {
|
|
||||||
return nil, errs.Wrap(err)
|
|
||||||
}
|
|
||||||
mm := make(map[string]int)
|
|
||||||
for k, v := range m {
|
|
||||||
state, err := strconv.Atoi(v)
|
|
||||||
if err != nil {
|
|
||||||
log.ZError(ctx, "token value is not int", err, "value", v, "userID", userID, "platformID", platformID)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
mm[strings.TrimPrefix(k, prefix)] = state
|
|
||||||
}
|
|
||||||
return mm, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *tokenCache) HasTemporaryToken(ctx context.Context, userID string, platformID int, token string) error {
|
|
||||||
key := cachekey.GetTemporaryTokenKey(userID, platformID, token)
|
|
||||||
if _, err := x.cache.Get(ctx, []string{key}); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *tokenCache) GetAllTokensWithoutError(ctx context.Context, userID string) (map[int]map[string]int, error) {
|
|
||||||
prefix := cachekey.UidPidToken + userID + ":"
|
|
||||||
tokens, err := x.cache.Prefix(ctx, prefix)
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
res := make(map[int]map[string]int)
|
|
||||||
for key, flagStr := range tokens {
|
|
||||||
flag, err := strconv.Atoi(flagStr)
|
|
||||||
if err != nil {
|
|
||||||
log.ZError(ctx, "token value is not int", err, "key", key, "value", flagStr, "userID", userID)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
arr := strings.SplitN(strings.TrimPrefix(key, prefix), ":", 2)
|
|
||||||
if len(arr) != 2 {
|
|
||||||
log.ZError(ctx, "token value is not int", err, "key", key, "value", flagStr, "userID", userID)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
platformID, err := strconv.Atoi(arr[0])
|
|
||||||
if err != nil {
|
|
||||||
log.ZError(ctx, "token value is not int", err, "key", key, "value", flagStr, "userID", userID)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
token := arr[1]
|
|
||||||
if token == "" {
|
|
||||||
log.ZError(ctx, "token value is not int", err, "key", key, "value", flagStr, "userID", userID)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
tk, ok := res[platformID]
|
|
||||||
if !ok {
|
|
||||||
tk = make(map[string]int)
|
|
||||||
res[platformID] = tk
|
|
||||||
}
|
|
||||||
tk[token] = flag
|
|
||||||
}
|
|
||||||
return res, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *tokenCache) SetTokenMapByUidPid(ctx context.Context, userID string, platformID int, m map[string]int) error {
|
|
||||||
for token, flag := range m {
|
|
||||||
err := x.SetTokenFlag(ctx, userID, platformID, token, flag)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *tokenCache) BatchSetTokenMapByUidPid(ctx context.Context, tokens map[string]map[string]any) error {
|
|
||||||
for prefix, tokenFlag := range tokens {
|
|
||||||
for token, flag := range tokenFlag {
|
|
||||||
flagStr := fmt.Sprintf("%v", flag)
|
|
||||||
if err := x.cache.Set(ctx, prefix+":"+token, flagStr, x.accessExpire); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *tokenCache) DeleteTokenByUidPid(ctx context.Context, userID string, platformID int, fields []string) error {
|
|
||||||
keys := make([]string, 0, len(fields))
|
|
||||||
for _, token := range fields {
|
|
||||||
keys = append(keys, x.getTokenKey(userID, platformID, token))
|
|
||||||
}
|
|
||||||
return x.cache.Del(ctx, keys)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *tokenCache) getExpireTime(t int64) time.Duration {
|
|
||||||
return time.Hour * 24 * time.Duration(t)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *tokenCache) DeleteTokenByTokenMap(ctx context.Context, userID string, tokens map[int][]string) error {
|
|
||||||
keys := make([]string, 0, len(tokens))
|
|
||||||
for platformID, ts := range tokens {
|
|
||||||
for _, t := range ts {
|
|
||||||
keys = append(keys, x.getTokenKey(userID, platformID, t))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return x.cache.Del(ctx, keys)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (x *tokenCache) DeleteAndSetTemporary(ctx context.Context, userID string, platformID int, fields []string) error {
|
|
||||||
keys := make([]string, 0, len(fields))
|
|
||||||
for _, f := range fields {
|
|
||||||
keys = append(keys, x.getTokenKey(userID, platformID, f))
|
|
||||||
}
|
|
||||||
if err := x.cache.Del(ctx, keys); err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
for _, f := range fields {
|
|
||||||
k := cachekey.GetTemporaryTokenKey(userID, platformID, f)
|
|
||||||
if err := x.cache.Set(ctx, k, "", time.Minute*5); err != nil {
|
|
||||||
return errs.Wrap(err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
@ -1,63 +0,0 @@
|
|||||||
package mcache
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"encoding/json"
|
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/database"
|
|
||||||
"github.com/openimsdk/tools/log"
|
|
||||||
)
|
|
||||||
|
|
||||||
func getCache[V any](ctx context.Context, cache database.Cache, key string, expireTime time.Duration, fn func(ctx context.Context) (V, error)) (V, error) {
|
|
||||||
getDB := func() (V, bool, error) {
|
|
||||||
res, err := cache.Get(ctx, []string{key})
|
|
||||||
if err != nil {
|
|
||||||
var val V
|
|
||||||
return val, false, err
|
|
||||||
}
|
|
||||||
var val V
|
|
||||||
if str, ok := res[key]; ok {
|
|
||||||
if json.Unmarshal([]byte(str), &val) != nil {
|
|
||||||
return val, false, err
|
|
||||||
}
|
|
||||||
return val, true, nil
|
|
||||||
}
|
|
||||||
return val, false, nil
|
|
||||||
}
|
|
||||||
dbVal, ok, err := getDB()
|
|
||||||
if err != nil {
|
|
||||||
return dbVal, err
|
|
||||||
}
|
|
||||||
if ok {
|
|
||||||
return dbVal, nil
|
|
||||||
}
|
|
||||||
lockValue, err := cache.Lock(ctx, key, time.Minute)
|
|
||||||
if err != nil {
|
|
||||||
return dbVal, err
|
|
||||||
}
|
|
||||||
defer func() {
|
|
||||||
if err := cache.Unlock(ctx, key, lockValue); err != nil {
|
|
||||||
log.ZError(ctx, "unlock cache key", err, "key", key, "value", lockValue)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
dbVal, ok, err = getDB()
|
|
||||||
if err != nil {
|
|
||||||
return dbVal, err
|
|
||||||
}
|
|
||||||
if ok {
|
|
||||||
return dbVal, nil
|
|
||||||
}
|
|
||||||
val, err := fn(ctx)
|
|
||||||
if err != nil {
|
|
||||||
return val, err
|
|
||||||
}
|
|
||||||
data, err := json.Marshal(val)
|
|
||||||
if err != nil {
|
|
||||||
return val, err
|
|
||||||
}
|
|
||||||
if err := cache.Set(ctx, key, string(data), expireTime); err != nil {
|
|
||||||
return val, err
|
|
||||||
}
|
|
||||||
return val, nil
|
|
||||||
}
|
|
||||||
@ -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
|
||||||
|
}
|
||||||
@ -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.ElementsMatch(t, []string{"127.0.0.1:10001", "127.0.0.1:10003"}, addrs)
|
||||||
|
assert.NoError(t, mock.ExpectationsWereMet())
|
||||||
|
}
|
||||||
Loading…
Reference in new issue