You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
409 lines
14 KiB
409 lines
14 KiB
1 year ago
|
// Copyright © 2023 OpenIM. All rights reserved.
|
||
|
//
|
||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||
|
// you may not use this file except in compliance with the License.
|
||
|
// You may obtain a copy of the License at
|
||
|
//
|
||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||
|
//
|
||
|
// Unless required by applicable law or agreed to in writing, software
|
||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||
|
// See the License for the specific language governing permissions and
|
||
|
// limitations under the License.
|
||
|
|
||
7 months ago
|
package redis
|
||
1 year ago
|
|
||
|
import (
|
||
|
"context"
|
||
1 year ago
|
"fmt"
|
||
1 year ago
|
"github.com/dtm-labs/rockscache"
|
||
9 months ago
|
"github.com/openimsdk/open-im-server/v3/pkg/common/config"
|
||
7 months ago
|
"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/common"
|
||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/database"
|
||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/storage/model"
|
||
8 months ago
|
"github.com/openimsdk/protocol/constant"
|
||
|
"github.com/openimsdk/tools/errs"
|
||
|
"github.com/openimsdk/tools/log"
|
||
|
"github.com/openimsdk/tools/utils/datautil"
|
||
9 months ago
|
"github.com/redis/go-redis/v9"
|
||
7 months ago
|
"time"
|
||
1 year ago
|
)
|
||
|
|
||
|
const (
|
||
9 months ago
|
groupExpireTime = time.Second * 60 * 60 * 12
|
||
1 year ago
|
)
|
||
|
|
||
7 months ago
|
var errIndex = errs.New("err index")
|
||
1 year ago
|
|
||
|
type GroupCacheRedis struct {
|
||
7 months ago
|
cache.BatchDeleter
|
||
|
groupDB database.Group
|
||
|
groupMemberDB database.GroupMember
|
||
|
groupRequestDB database.GroupRequest
|
||
1 year ago
|
expireTime time.Duration
|
||
|
rcClient *rockscache.Client
|
||
7 months ago
|
groupHash cache.GroupHash
|
||
1 year ago
|
}
|
||
|
|
||
1 year ago
|
func NewGroupCacheRedis(
|
||
|
rdb redis.UniversalClient,
|
||
8 months ago
|
localCache *config.LocalCache,
|
||
7 months ago
|
groupDB database.Group,
|
||
|
groupMemberDB database.GroupMember,
|
||
|
groupRequestDB database.GroupRequest,
|
||
|
hashCode cache.GroupHash,
|
||
|
opts *rockscache.Options,
|
||
|
) cache.GroupCache {
|
||
|
batchHandler := NewBatchDeleterRedis(rdb, opts, []string{localCache.Group.Topic})
|
||
8 months ago
|
g := localCache.Group
|
||
9 months ago
|
log.ZDebug(context.Background(), "group local cache init", "Topic", g.Topic, "SlotNum", g.SlotNum, "SlotSize", g.SlotSize, "enable", g.Enable())
|
||
7 months ago
|
|
||
1 year ago
|
return &GroupCacheRedis{
|
||
7 months ago
|
BatchDeleter: batchHandler,
|
||
|
rcClient: rockscache.NewClient(rdb, *opts),
|
||
|
expireTime: groupExpireTime,
|
||
|
groupDB: groupDB,
|
||
|
groupMemberDB: groupMemberDB,
|
||
|
groupRequestDB: groupRequestDB,
|
||
|
groupHash: hashCode,
|
||
1 year ago
|
}
|
||
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) CloneGroupCache() cache.GroupCache {
|
||
1 year ago
|
return &GroupCacheRedis{
|
||
7 months ago
|
BatchDeleter: g.BatchDeleter.Clone(),
|
||
1 year ago
|
rcClient: g.rcClient,
|
||
|
expireTime: g.expireTime,
|
||
|
groupDB: g.groupDB,
|
||
|
groupMemberDB: g.groupMemberDB,
|
||
|
groupRequestDB: g.groupRequestDB,
|
||
|
}
|
||
1 year ago
|
}
|
||
|
|
||
|
func (g *GroupCacheRedis) getGroupInfoKey(groupID string) string {
|
||
9 months ago
|
return cachekey.GetGroupInfoKey(groupID)
|
||
1 year ago
|
}
|
||
|
|
||
|
func (g *GroupCacheRedis) getJoinedGroupsKey(userID string) string {
|
||
9 months ago
|
return cachekey.GetJoinedGroupsKey(userID)
|
||
1 year ago
|
}
|
||
|
|
||
|
func (g *GroupCacheRedis) getGroupMembersHashKey(groupID string) string {
|
||
9 months ago
|
return cachekey.GetGroupMembersHashKey(groupID)
|
||
1 year ago
|
}
|
||
|
|
||
|
func (g *GroupCacheRedis) getGroupMemberIDsKey(groupID string) string {
|
||
9 months ago
|
return cachekey.GetGroupMemberIDsKey(groupID)
|
||
1 year ago
|
}
|
||
|
|
||
|
func (g *GroupCacheRedis) getGroupMemberInfoKey(groupID, userID string) string {
|
||
9 months ago
|
return cachekey.GetGroupMemberInfoKey(groupID, userID)
|
||
1 year ago
|
}
|
||
|
|
||
|
func (g *GroupCacheRedis) getGroupMemberNumKey(groupID string) string {
|
||
9 months ago
|
return cachekey.GetGroupMemberNumKey(groupID)
|
||
1 year ago
|
}
|
||
|
|
||
1 year ago
|
func (g *GroupCacheRedis) getGroupRoleLevelMemberIDsKey(groupID string, roleLevel int32) string {
|
||
9 months ago
|
return cachekey.GetGroupRoleLevelMemberIDsKey(groupID, roleLevel)
|
||
1 year ago
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetGroupIndex(group *model.Group, keys []string) (int, error) {
|
||
1 year ago
|
key := g.getGroupInfoKey(group.GroupID)
|
||
|
for i, _key := range keys {
|
||
|
if _key == key {
|
||
|
return i, nil
|
||
|
}
|
||
|
}
|
||
1 year ago
|
|
||
1 year ago
|
return 0, errIndex
|
||
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetGroupMemberIndex(groupMember *model.GroupMember, keys []string) (int, error) {
|
||
1 year ago
|
key := g.getGroupMemberInfoKey(groupMember.GroupID, groupMember.UserID)
|
||
|
for i, _key := range keys {
|
||
|
if _key == key {
|
||
|
return i, nil
|
||
|
}
|
||
|
}
|
||
1 year ago
|
|
||
1 year ago
|
return 0, errIndex
|
||
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetGroupsInfo(ctx context.Context, groupIDs []string) (groups []*model.Group, err error) {
|
||
|
return batchGetCache(ctx, g.rcClient, g.expireTime, groupIDs, func(groupID string) string {
|
||
1 year ago
|
return g.getGroupInfoKey(groupID)
|
||
7 months ago
|
}, func(ctx context.Context, groupID string) (*model.Group, error) {
|
||
1 year ago
|
return g.groupDB.Take(ctx, groupID)
|
||
|
})
|
||
1 year ago
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetGroupInfo(ctx context.Context, groupID string) (group *model.Group, err error) {
|
||
|
return getCache(ctx, g.rcClient, g.getGroupInfoKey(groupID), g.expireTime, func(ctx context.Context) (*model.Group, error) {
|
||
1 year ago
|
return g.groupDB.Take(ctx, groupID)
|
||
|
})
|
||
1 year ago
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) DelGroupsInfo(groupIDs ...string) cache.GroupCache {
|
||
|
newGroupCache := g.CloneGroupCache()
|
||
1 year ago
|
keys := make([]string, 0, len(groupIDs))
|
||
1 year ago
|
for _, groupID := range groupIDs {
|
||
|
keys = append(keys, g.getGroupInfoKey(groupID))
|
||
|
}
|
||
1 year ago
|
newGroupCache.AddKeys(keys...)
|
||
|
|
||
|
return newGroupCache
|
||
1 year ago
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) DelGroupsOwner(groupIDs ...string) cache.GroupCache {
|
||
|
newGroupCache := g.CloneGroupCache()
|
||
1 year ago
|
keys := make([]string, 0, len(groupIDs))
|
||
|
for _, groupID := range groupIDs {
|
||
|
keys = append(keys, g.getGroupRoleLevelMemberIDsKey(groupID, constant.GroupOwner))
|
||
1 year ago
|
}
|
||
1 year ago
|
newGroupCache.AddKeys(keys...)
|
||
|
|
||
|
return newGroupCache
|
||
1 year ago
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) DelGroupRoleLevel(groupID string, roleLevels []int32) cache.GroupCache {
|
||
|
newGroupCache := g.CloneGroupCache()
|
||
1 year ago
|
keys := make([]string, 0, len(roleLevels))
|
||
|
for _, roleLevel := range roleLevels {
|
||
|
keys = append(keys, g.getGroupRoleLevelMemberIDsKey(groupID, roleLevel))
|
||
1 year ago
|
}
|
||
1 year ago
|
newGroupCache.AddKeys(keys...)
|
||
|
return newGroupCache
|
||
1 year ago
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) DelGroupAllRoleLevel(groupID string) cache.GroupCache {
|
||
1 year ago
|
return g.DelGroupRoleLevel(groupID, []int32{constant.GroupOwner, constant.GroupAdmin, constant.GroupOrdinaryUsers})
|
||
|
}
|
||
|
|
||
1 year ago
|
func (g *GroupCacheRedis) GetGroupMembersHash(ctx context.Context, groupID string) (hashCode uint64, err error) {
|
||
1 year ago
|
if g.groupHash == nil {
|
||
8 months ago
|
return 0, errs.ErrInternalServer.WrapMsg("group hash is nil")
|
||
1 year ago
|
}
|
||
1 year ago
|
return getCache(ctx, g.rcClient, g.getGroupMembersHashKey(groupID), g.expireTime, func(ctx context.Context) (uint64, error) {
|
||
1 year ago
|
return g.groupHash.GetGroupHash(ctx, groupID)
|
||
1 year ago
|
})
|
||
1 year ago
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetGroupMemberHashMap(ctx context.Context, groupIDs []string) (map[string]*common.GroupSimpleUserID, error) {
|
||
1 year ago
|
if g.groupHash == nil {
|
||
8 months ago
|
return nil, errs.ErrInternalServer.WrapMsg("group hash is nil")
|
||
1 year ago
|
}
|
||
7 months ago
|
res := make(map[string]*common.GroupSimpleUserID)
|
||
1 year ago
|
for _, groupID := range groupIDs {
|
||
|
hash, err := g.GetGroupMembersHash(ctx, groupID)
|
||
|
if err != nil {
|
||
|
return nil, err
|
||
|
}
|
||
8 months ago
|
log.ZDebug(ctx, "GetGroupMemberHashMap", "groupID", groupID, "hash", hash)
|
||
1 year ago
|
num, err := g.GetGroupMemberNum(ctx, groupID)
|
||
|
if err != nil {
|
||
|
return nil, err
|
||
|
}
|
||
7 months ago
|
res[groupID] = &common.GroupSimpleUserID{Hash: hash, MemberNum: uint32(num)}
|
||
1 year ago
|
}
|
||
1 year ago
|
|
||
1 year ago
|
return res, nil
|
||
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) DelGroupMembersHash(groupID string) cache.GroupCache {
|
||
|
cache := g.CloneGroupCache()
|
||
1 year ago
|
cache.AddKeys(g.getGroupMembersHashKey(groupID))
|
||
1 year ago
|
|
||
1 year ago
|
return cache
|
||
|
}
|
||
|
|
||
|
func (g *GroupCacheRedis) GetGroupMemberIDs(ctx context.Context, groupID string) (groupMemberIDs []string, err error) {
|
||
1 year ago
|
return getCache(ctx, g.rcClient, g.getGroupMemberIDsKey(groupID), g.expireTime, func(ctx context.Context) ([]string, error) {
|
||
|
return g.groupMemberDB.FindMemberUserID(ctx, groupID)
|
||
|
})
|
||
1 year ago
|
}
|
||
|
|
||
|
func (g *GroupCacheRedis) GetGroupsMemberIDs(ctx context.Context, groupIDs []string) (map[string][]string, error) {
|
||
|
m := make(map[string][]string)
|
||
|
for _, groupID := range groupIDs {
|
||
|
userIDs, err := g.GetGroupMemberIDs(ctx, groupID)
|
||
|
if err != nil {
|
||
|
return nil, err
|
||
|
}
|
||
|
m[groupID] = userIDs
|
||
|
}
|
||
1 year ago
|
|
||
1 year ago
|
return m, nil
|
||
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) DelGroupMemberIDs(groupID string) cache.GroupCache {
|
||
|
cache := g.CloneGroupCache()
|
||
1 year ago
|
cache.AddKeys(g.getGroupMemberIDsKey(groupID))
|
||
1 year ago
|
|
||
1 year ago
|
return cache
|
||
|
}
|
||
|
|
||
|
func (g *GroupCacheRedis) GetJoinedGroupIDs(ctx context.Context, userID string) (joinedGroupIDs []string, err error) {
|
||
1 year ago
|
return getCache(ctx, g.rcClient, g.getJoinedGroupsKey(userID), g.expireTime, func(ctx context.Context) ([]string, error) {
|
||
|
return g.groupMemberDB.FindUserJoinedGroupID(ctx, userID)
|
||
|
})
|
||
1 year ago
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) DelJoinedGroupID(userIDs ...string) cache.GroupCache {
|
||
1 year ago
|
keys := make([]string, 0, len(userIDs))
|
||
1 year ago
|
for _, userID := range userIDs {
|
||
|
keys = append(keys, g.getJoinedGroupsKey(userID))
|
||
|
}
|
||
7 months ago
|
cache := g.CloneGroupCache()
|
||
1 year ago
|
cache.AddKeys(keys...)
|
||
1 year ago
|
|
||
1 year ago
|
return cache
|
||
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetGroupMemberInfo(ctx context.Context, groupID, userID string) (groupMember *model.GroupMember, err error) {
|
||
|
return getCache(ctx, g.rcClient, g.getGroupMemberInfoKey(groupID, userID), g.expireTime, func(ctx context.Context) (*model.GroupMember, error) {
|
||
1 year ago
|
return g.groupMemberDB.Take(ctx, groupID, userID)
|
||
|
})
|
||
1 year ago
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetGroupMembersInfo(ctx context.Context, groupID string, userIDs []string) ([]*model.GroupMember, error) {
|
||
|
return batchGetCache(ctx, g.rcClient, g.expireTime, userIDs, func(userID string) string {
|
||
1 year ago
|
return g.getGroupMemberInfoKey(groupID, userID)
|
||
7 months ago
|
}, func(ctx context.Context, userID string) (*model.GroupMember, error) {
|
||
1 year ago
|
return g.groupMemberDB.Take(ctx, groupID, userID)
|
||
|
})
|
||
1 year ago
|
}
|
||
|
|
||
1 year ago
|
func (g *GroupCacheRedis) GetGroupMembersPage(
|
||
|
ctx context.Context,
|
||
|
groupID string,
|
||
|
userIDs []string,
|
||
|
showNumber, pageNumber int32,
|
||
7 months ago
|
) (total uint32, groupMembers []*model.GroupMember, err error) {
|
||
1 year ago
|
groupMemberIDs, err := g.GetGroupMemberIDs(ctx, groupID)
|
||
|
if err != nil {
|
||
|
return 0, nil, err
|
||
|
}
|
||
|
if userIDs != nil {
|
||
8 months ago
|
userIDs = datautil.BothExist(userIDs, groupMemberIDs)
|
||
1 year ago
|
} else {
|
||
|
userIDs = groupMemberIDs
|
||
|
}
|
||
8 months ago
|
groupMembers, err = g.GetGroupMembersInfo(ctx, groupID, datautil.Paginate(userIDs, int(showNumber), int(showNumber)))
|
||
1 year ago
|
|
||
1 year ago
|
return uint32(len(userIDs)), groupMembers, err
|
||
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetAllGroupMembersInfo(ctx context.Context, groupID string) (groupMembers []*model.GroupMember, err error) {
|
||
1 year ago
|
groupMemberIDs, err := g.GetGroupMemberIDs(ctx, groupID)
|
||
|
if err != nil {
|
||
|
return nil, err
|
||
|
}
|
||
1 year ago
|
|
||
1 year ago
|
return g.GetGroupMembersInfo(ctx, groupID, groupMemberIDs)
|
||
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetAllGroupMemberInfo(ctx context.Context, groupID string) ([]*model.GroupMember, error) {
|
||
1 year ago
|
groupMemberIDs, err := g.GetGroupMemberIDs(ctx, groupID)
|
||
|
if err != nil {
|
||
|
return nil, err
|
||
|
}
|
||
1 year ago
|
return g.GetGroupMembersInfo(ctx, groupID, groupMemberIDs)
|
||
1 year ago
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) DelGroupMembersInfo(groupID string, userIDs ...string) cache.GroupCache {
|
||
1 year ago
|
keys := make([]string, 0, len(userIDs))
|
||
1 year ago
|
for _, userID := range userIDs {
|
||
|
keys = append(keys, g.getGroupMemberInfoKey(groupID, userID))
|
||
|
}
|
||
7 months ago
|
cache := g.CloneGroupCache()
|
||
1 year ago
|
cache.AddKeys(keys...)
|
||
1 year ago
|
|
||
1 year ago
|
return cache
|
||
|
}
|
||
|
|
||
|
func (g *GroupCacheRedis) GetGroupMemberNum(ctx context.Context, groupID string) (memberNum int64, err error) {
|
||
1 year ago
|
return getCache(ctx, g.rcClient, g.getGroupMemberNumKey(groupID), g.expireTime, func(ctx context.Context) (int64, error) {
|
||
|
return g.groupMemberDB.TakeGroupMemberNum(ctx, groupID)
|
||
|
})
|
||
1 year ago
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) DelGroupsMemberNum(groupID ...string) cache.GroupCache {
|
||
1 year ago
|
keys := make([]string, 0, len(groupID))
|
||
1 year ago
|
for _, groupID := range groupID {
|
||
|
keys = append(keys, g.getGroupMemberNumKey(groupID))
|
||
|
}
|
||
7 months ago
|
cache := g.CloneGroupCache()
|
||
1 year ago
|
cache.AddKeys(keys...)
|
||
1 year ago
|
|
||
1 year ago
|
return cache
|
||
|
}
|
||
1 year ago
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetGroupOwner(ctx context.Context, groupID string) (*model.GroupMember, error) {
|
||
1 year ago
|
members, err := g.GetGroupRoleLevelMemberInfo(ctx, groupID, constant.GroupOwner)
|
||
|
if err != nil {
|
||
|
return nil, err
|
||
|
}
|
||
|
if len(members) == 0 {
|
||
8 months ago
|
return nil, errs.ErrRecordNotFound.WrapMsg(fmt.Sprintf("group %s owner not found", groupID))
|
||
1 year ago
|
}
|
||
|
return members[0], nil
|
||
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetGroupsOwner(ctx context.Context, groupIDs []string) ([]*model.GroupMember, error) {
|
||
|
members := make([]*model.GroupMember, 0, len(groupIDs))
|
||
1 year ago
|
for _, groupID := range groupIDs {
|
||
|
items, err := g.GetGroupRoleLevelMemberInfo(ctx, groupID, constant.GroupOwner)
|
||
|
if err != nil {
|
||
|
return nil, err
|
||
|
}
|
||
|
if len(items) > 0 {
|
||
|
members = append(members, items[0])
|
||
|
}
|
||
|
}
|
||
|
return members, nil
|
||
|
}
|
||
|
|
||
|
func (g *GroupCacheRedis) GetGroupRoleLevelMemberIDs(ctx context.Context, groupID string, roleLevel int32) ([]string, error) {
|
||
|
return getCache(ctx, g.rcClient, g.getGroupRoleLevelMemberIDsKey(groupID, roleLevel), g.expireTime, func(ctx context.Context) ([]string, error) {
|
||
|
return g.groupMemberDB.FindRoleLevelUserIDs(ctx, groupID, roleLevel)
|
||
|
})
|
||
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetGroupRoleLevelMemberInfo(ctx context.Context, groupID string, roleLevel int32) ([]*model.GroupMember, error) {
|
||
1 year ago
|
userIDs, err := g.GetGroupRoleLevelMemberIDs(ctx, groupID, roleLevel)
|
||
|
if err != nil {
|
||
|
return nil, err
|
||
|
}
|
||
|
return g.GetGroupMembersInfo(ctx, groupID, userIDs)
|
||
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) GetGroupRolesLevelMemberInfo(ctx context.Context, groupID string, roleLevels []int32) ([]*model.GroupMember, error) {
|
||
1 year ago
|
var userIDs []string
|
||
|
for _, roleLevel := range roleLevels {
|
||
|
ids, err := g.GetGroupRoleLevelMemberIDs(ctx, groupID, roleLevel)
|
||
|
if err != nil {
|
||
|
return nil, err
|
||
|
}
|
||
|
userIDs = append(userIDs, ids...)
|
||
|
}
|
||
|
return g.GetGroupMembersInfo(ctx, groupID, userIDs)
|
||
|
}
|
||
|
|
||
7 months ago
|
func (g *GroupCacheRedis) FindGroupMemberUser(ctx context.Context, groupIDs []string, userID string) (_ []*model.GroupMember, err error) {
|
||
1 year ago
|
if len(groupIDs) == 0 {
|
||
|
groupIDs, err = g.GetJoinedGroupIDs(ctx, userID)
|
||
|
if err != nil {
|
||
|
return nil, err
|
||
|
}
|
||
|
}
|
||
7 months ago
|
return batchGetCache(ctx, g.rcClient, g.expireTime, groupIDs, func(groupID string) string {
|
||
1 year ago
|
return g.getGroupMemberInfoKey(groupID, userID)
|
||
7 months ago
|
}, func(ctx context.Context, groupID string) (*model.GroupMember, error) {
|
||
1 year ago
|
return g.groupMemberDB.Take(ctx, groupID, userID)
|
||
|
})
|
||
|
}
|