feat: add local cache for high frequency reads (#2036)
* feat: msg local cache * feat: msg local cache * feat: msg local cache * feat: msg local cache * feat: msg local cache * feat: msg local cache * fix: mongo * fix: mongo * fix: mongo * openim.yaml * localcache * localcache * localcache * localcache * localcache * localcache * localcache * localcache * localcache * local cache * local cache * local cache * local cache * fix: GroupApplicationAcceptedNotification * fix: GroupApplicationAcceptedNotification * fix: NotificationUserInfoUpdate * feat: cache add single-flight and timing-wheel. * feat: local cache * feat: local cache * feat: local cache * feat: cache add single-flight and timing-wheel. * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: local cache * feat: msg rpc local cache * feat: msg rpc local cache * feat: msg rpc local cache * feat: msg rpc local cache * feat: msg rpc local cache * feat: msg rpc local cache * refactor: refactor the code of push and optimization. * cicd: robot automated Change * refactor: rename cache. * merge * fix: refactor project dir avoid import cycle. * update tools * merge * feat: conversation FindRecvMsgNotNotifyUserIDs * feat: conversation FindRecvMsgNotNotifyUserIDs * feat: conversation FindRecvMsgNotNotifyUserIDs * merge * merge the latest main --------- Co-authored-by: Gordon <46924906+FGadvancer@users.noreply.github.com> Co-authored-by: withchao <withchao@users.noreply.github.com>pull/2047/head
parent
291443dd6b
commit
b9cf40034c
@ -0,0 +1,15 @@
|
|||||||
|
package cachekey
|
||||||
|
|
||||||
|
const (
|
||||||
|
BlackIDsKey = "BLACK_IDS:"
|
||||||
|
IsBlackKey = "IS_BLACK:" // local cache
|
||||||
|
)
|
||||||
|
|
||||||
|
func GetBlackIDsKey(ownerUserID string) string {
|
||||||
|
return BlackIDsKey + ownerUserID
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetIsBlackIDsKey(possibleBlackUserID, userID string) string {
|
||||||
|
return IsBlackKey + userID + "-" + possibleBlackUserID
|
||||||
|
}
|
@ -0,0 +1,44 @@
|
|||||||
|
package cachekey
|
||||||
|
|
||||||
|
const (
|
||||||
|
ConversationKey = "CONVERSATION:"
|
||||||
|
ConversationIDsKey = "CONVERSATION_IDS:"
|
||||||
|
ConversationIDsHashKey = "CONVERSATION_IDS_HASH:"
|
||||||
|
ConversationHasReadSeqKey = "CONVERSATION_HAS_READ_SEQ:"
|
||||||
|
RecvMsgOptKey = "RECV_MSG_OPT:"
|
||||||
|
SuperGroupRecvMsgNotNotifyUserIDsKey = "SUPER_GROUP_RECV_MSG_NOT_NOTIFY_USER_IDS:"
|
||||||
|
SuperGroupRecvMsgNotNotifyUserIDsHashKey = "SUPER_GROUP_RECV_MSG_NOT_NOTIFY_USER_IDS_HASH:"
|
||||||
|
ConversationNotReceiveMessageUserIDsKey = "CONVERSATION_NOT_RECEIVE_MESSAGE_USER_IDS:"
|
||||||
|
)
|
||||||
|
|
||||||
|
func GetConversationKey(ownerUserID, conversationID string) string {
|
||||||
|
return ConversationKey + ownerUserID + ":" + conversationID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetConversationIDsKey(ownerUserID string) string {
|
||||||
|
return ConversationIDsKey + ownerUserID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetSuperGroupRecvNotNotifyUserIDsKey(groupID string) string {
|
||||||
|
return SuperGroupRecvMsgNotNotifyUserIDsKey + groupID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetRecvMsgOptKey(ownerUserID, conversationID string) string {
|
||||||
|
return RecvMsgOptKey + ownerUserID + ":" + conversationID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetSuperGroupRecvNotNotifyUserIDsHashKey(groupID string) string {
|
||||||
|
return SuperGroupRecvMsgNotNotifyUserIDsHashKey + groupID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetConversationHasReadSeqKey(ownerUserID, conversationID string) string {
|
||||||
|
return ConversationHasReadSeqKey + ownerUserID + ":" + conversationID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetConversationNotReceiveMessageUserIDsKey(conversationID string) string {
|
||||||
|
return ConversationNotReceiveMessageUserIDsKey + conversationID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetUserConversationIDsHashKey(ownerUserID string) string {
|
||||||
|
return ConversationIDsHashKey + ownerUserID
|
||||||
|
}
|
@ -0,0 +1,24 @@
|
|||||||
|
package cachekey
|
||||||
|
|
||||||
|
const (
|
||||||
|
FriendIDsKey = "FRIEND_IDS:"
|
||||||
|
TwoWayFriendsIDsKey = "COMMON_FRIENDS_IDS:"
|
||||||
|
FriendKey = "FRIEND_INFO:"
|
||||||
|
IsFriendKey = "IS_FRIEND:" // local cache key
|
||||||
|
)
|
||||||
|
|
||||||
|
func GetFriendIDsKey(ownerUserID string) string {
|
||||||
|
return FriendIDsKey + ownerUserID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetTwoWayFriendsIDsKey(ownerUserID string) string {
|
||||||
|
return TwoWayFriendsIDsKey + ownerUserID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetFriendKey(ownerUserID, friendUserID string) string {
|
||||||
|
return FriendKey + ownerUserID + "-" + friendUserID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetIsFriendKey(possibleFriendUserID, userID string) string {
|
||||||
|
return IsFriendKey + possibleFriendUserID + "-" + userID
|
||||||
|
}
|
@ -0,0 +1,45 @@
|
|||||||
|
package cachekey
|
||||||
|
|
||||||
|
import (
|
||||||
|
"strconv"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
groupExpireTime = time.Second * 60 * 60 * 12
|
||||||
|
GroupInfoKey = "GROUP_INFO:"
|
||||||
|
GroupMemberIDsKey = "GROUP_MEMBER_IDS:"
|
||||||
|
GroupMembersHashKey = "GROUP_MEMBERS_HASH2:"
|
||||||
|
GroupMemberInfoKey = "GROUP_MEMBER_INFO:"
|
||||||
|
JoinedGroupsKey = "JOIN_GROUPS_KEY:"
|
||||||
|
GroupMemberNumKey = "GROUP_MEMBER_NUM_CACHE:"
|
||||||
|
GroupRoleLevelMemberIDsKey = "GROUP_ROLE_LEVEL_MEMBER_IDS:"
|
||||||
|
)
|
||||||
|
|
||||||
|
func GetGroupInfoKey(groupID string) string {
|
||||||
|
return GroupInfoKey + groupID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetJoinedGroupsKey(userID string) string {
|
||||||
|
return JoinedGroupsKey + userID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetGroupMembersHashKey(groupID string) string {
|
||||||
|
return GroupMembersHashKey + groupID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetGroupMemberIDsKey(groupID string) string {
|
||||||
|
return GroupMemberIDsKey + groupID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetGroupMemberInfoKey(groupID, userID string) string {
|
||||||
|
return GroupMemberInfoKey + groupID + "-" + userID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetGroupMemberNumKey(groupID string) string {
|
||||||
|
return GroupMemberNumKey + groupID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetGroupRoleLevelMemberIDsKey(groupID string, roleLevel int32) string {
|
||||||
|
return GroupRoleLevelMemberIDsKey + groupID + "-" + strconv.Itoa(int(roleLevel))
|
||||||
|
}
|
@ -0,0 +1,14 @@
|
|||||||
|
package cachekey
|
||||||
|
|
||||||
|
const (
|
||||||
|
UserInfoKey = "USER_INFO:"
|
||||||
|
UserGlobalRecvMsgOptKey = "USER_GLOBAL_RECV_MSG_OPT_KEY:"
|
||||||
|
)
|
||||||
|
|
||||||
|
func GetUserInfoKey(userID string) string {
|
||||||
|
return UserInfoKey + userID
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetUserGlobalRecvMsgOptKey(userID string) string {
|
||||||
|
return UserGlobalRecvMsgOptKey + userID
|
||||||
|
}
|
@ -0,0 +1,66 @@
|
|||||||
|
package cache
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/cachekey"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/config"
|
||||||
|
"strings"
|
||||||
|
"sync"
|
||||||
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
once sync.Once
|
||||||
|
subscribe map[string][]string
|
||||||
|
)
|
||||||
|
|
||||||
|
func getPublishKey(topic string, key []string) []string {
|
||||||
|
if topic == "" || len(key) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
once.Do(func() {
|
||||||
|
list := []struct {
|
||||||
|
Local config.LocalCache
|
||||||
|
Keys []string
|
||||||
|
}{
|
||||||
|
{
|
||||||
|
Local: config.Config.LocalCache.User,
|
||||||
|
Keys: []string{cachekey.UserInfoKey, cachekey.UserGlobalRecvMsgOptKey},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
Local: config.Config.LocalCache.Group,
|
||||||
|
Keys: []string{cachekey.GroupMemberIDsKey, cachekey.GroupInfoKey, cachekey.GroupMemberInfoKey},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
Local: config.Config.LocalCache.Friend,
|
||||||
|
Keys: []string{cachekey.FriendIDsKey, cachekey.BlackIDsKey},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
Local: config.Config.LocalCache.Conversation,
|
||||||
|
Keys: []string{cachekey.ConversationKey, cachekey.ConversationIDsKey, cachekey.ConversationNotReceiveMessageUserIDsKey},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
subscribe = make(map[string][]string)
|
||||||
|
for _, v := range list {
|
||||||
|
if v.Local.Enable() {
|
||||||
|
subscribe[v.Local.Topic] = v.Keys
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
prefix, ok := subscribe[topic]
|
||||||
|
if !ok {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
res := make([]string, 0, len(key))
|
||||||
|
for _, k := range key {
|
||||||
|
var exist bool
|
||||||
|
for _, p := range prefix {
|
||||||
|
if strings.HasPrefix(k, p) {
|
||||||
|
exist = true
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if exist {
|
||||||
|
res = append(res, k)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return res
|
||||||
|
}
|
@ -1,86 +0,0 @@
|
|||||||
// 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.
|
|
||||||
|
|
||||||
package localcache
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"sync"
|
|
||||||
|
|
||||||
"github.com/OpenIMSDK/protocol/conversation"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/rpcclient"
|
|
||||||
)
|
|
||||||
|
|
||||||
type ConversationLocalCache struct {
|
|
||||||
lock sync.Mutex
|
|
||||||
superGroupRecvMsgNotNotifyUserIDs map[string]Hash
|
|
||||||
conversationIDs map[string]Hash
|
|
||||||
client *rpcclient.ConversationRpcClient
|
|
||||||
}
|
|
||||||
|
|
||||||
type Hash struct {
|
|
||||||
hash uint64
|
|
||||||
ids []string
|
|
||||||
}
|
|
||||||
|
|
||||||
func NewConversationLocalCache(client *rpcclient.ConversationRpcClient) *ConversationLocalCache {
|
|
||||||
return &ConversationLocalCache{
|
|
||||||
superGroupRecvMsgNotNotifyUserIDs: make(map[string]Hash),
|
|
||||||
conversationIDs: make(map[string]Hash),
|
|
||||||
client: client,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (g *ConversationLocalCache) GetRecvMsgNotNotifyUserIDs(ctx context.Context, groupID string) ([]string, error) {
|
|
||||||
resp, err := g.client.Client.GetRecvMsgNotNotifyUserIDs(ctx, &conversation.GetRecvMsgNotNotifyUserIDsReq{
|
|
||||||
GroupID: groupID,
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
return resp.UserIDs, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (g *ConversationLocalCache) GetConversationIDs(ctx context.Context, userID string) ([]string, error) {
|
|
||||||
resp, err := g.client.Client.GetUserConversationIDsHash(ctx, &conversation.GetUserConversationIDsHashReq{
|
|
||||||
OwnerUserID: userID,
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
|
|
||||||
g.lock.Lock()
|
|
||||||
hash, ok := g.conversationIDs[userID]
|
|
||||||
g.lock.Unlock()
|
|
||||||
|
|
||||||
if !ok || hash.hash != resp.Hash {
|
|
||||||
conversationIDsResp, err := g.client.Client.GetConversationIDs(ctx, &conversation.GetConversationIDsReq{
|
|
||||||
UserID: userID,
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
|
|
||||||
g.lock.Lock()
|
|
||||||
defer g.lock.Unlock()
|
|
||||||
g.conversationIDs[userID] = Hash{
|
|
||||||
hash: resp.Hash,
|
|
||||||
ids: conversationIDsResp.ConversationIDs,
|
|
||||||
}
|
|
||||||
|
|
||||||
return conversationIDsResp.ConversationIDs, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
return hash.ids, nil
|
|
||||||
}
|
|
@ -1,15 +0,0 @@
|
|||||||
// 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.
|
|
||||||
|
|
||||||
package localcache // import "github.com/openimsdk/open-im-server/v3/pkg/common/db/localcache"
|
|
@ -1,77 +0,0 @@
|
|||||||
// 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.
|
|
||||||
|
|
||||||
package localcache
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"sync"
|
|
||||||
|
|
||||||
"github.com/OpenIMSDK/protocol/group"
|
|
||||||
"github.com/OpenIMSDK/tools/errs"
|
|
||||||
"github.com/openimsdk/open-im-server/v3/pkg/rpcclient"
|
|
||||||
)
|
|
||||||
|
|
||||||
type GroupLocalCache struct {
|
|
||||||
lock sync.Mutex
|
|
||||||
cache map[string]GroupMemberIDsHash
|
|
||||||
client *rpcclient.GroupRpcClient
|
|
||||||
}
|
|
||||||
|
|
||||||
type GroupMemberIDsHash struct {
|
|
||||||
memberListHash uint64
|
|
||||||
userIDs []string
|
|
||||||
}
|
|
||||||
|
|
||||||
func NewGroupLocalCache(client *rpcclient.GroupRpcClient) *GroupLocalCache {
|
|
||||||
return &GroupLocalCache{
|
|
||||||
cache: make(map[string]GroupMemberIDsHash, 0),
|
|
||||||
client: client,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (g *GroupLocalCache) GetGroupMemberIDs(ctx context.Context, groupID string) ([]string, error) {
|
|
||||||
resp, err := g.client.Client.GetGroupAbstractInfo(ctx, &group.GetGroupAbstractInfoReq{
|
|
||||||
GroupIDs: []string{groupID},
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
if len(resp.GroupAbstractInfos) < 1 {
|
|
||||||
return nil, errs.ErrGroupIDNotFound
|
|
||||||
}
|
|
||||||
|
|
||||||
g.lock.Lock()
|
|
||||||
localHashInfo, ok := g.cache[groupID]
|
|
||||||
if ok && localHashInfo.memberListHash == resp.GroupAbstractInfos[0].GroupMemberListHash {
|
|
||||||
g.lock.Unlock()
|
|
||||||
return localHashInfo.userIDs, nil
|
|
||||||
}
|
|
||||||
g.lock.Unlock()
|
|
||||||
|
|
||||||
groupMembersResp, err := g.client.Client.GetGroupMemberUserIDs(ctx, &group.GetGroupMemberUserIDsReq{
|
|
||||||
GroupID: groupID,
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
|
|
||||||
g.lock.Lock()
|
|
||||||
defer g.lock.Unlock()
|
|
||||||
g.cache[groupID] = GroupMemberIDsHash{
|
|
||||||
memberListHash: resp.GroupAbstractInfos[0].GroupMemberListHash,
|
|
||||||
userIDs: groupMembersResp.UserIDs,
|
|
||||||
}
|
|
||||||
return g.cache[groupID].userIDs, nil
|
|
||||||
}
|
|
@ -1,15 +0,0 @@
|
|||||||
// 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.
|
|
||||||
|
|
||||||
package localcache
|
|
@ -0,0 +1,16 @@
|
|||||||
|
package redispubsub
|
||||||
|
|
||||||
|
import "github.com/redis/go-redis/v9"
|
||||||
|
|
||||||
|
type Publisher struct {
|
||||||
|
client redis.UniversalClient
|
||||||
|
channel string
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewPublisher(client redis.UniversalClient, channel string) *Publisher {
|
||||||
|
return &Publisher{client: client, channel: channel}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (p *Publisher) Publish(message string) error {
|
||||||
|
return p.client.Publish(ctx, p.channel, message).Err()
|
||||||
|
}
|
@ -0,0 +1,34 @@
|
|||||||
|
package redispubsub
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"github.com/redis/go-redis/v9"
|
||||||
|
)
|
||||||
|
|
||||||
|
var ctx = context.Background()
|
||||||
|
|
||||||
|
type Subscriber struct {
|
||||||
|
client redis.UniversalClient
|
||||||
|
channel string
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewSubscriber(client redis.UniversalClient, channel string) *Subscriber {
|
||||||
|
return &Subscriber{client: client, channel: channel}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Subscriber) OnMessage(ctx context.Context, callback func(string)) error {
|
||||||
|
messageChannel := s.client.Subscribe(ctx, s.channel).Channel()
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return
|
||||||
|
case msg := <-messageChannel:
|
||||||
|
callback(msg.Payload)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
@ -0,0 +1,112 @@
|
|||||||
|
package localcache
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"github.com/openimsdk/localcache/link"
|
||||||
|
"github.com/openimsdk/localcache/lru"
|
||||||
|
"hash/fnv"
|
||||||
|
"unsafe"
|
||||||
|
)
|
||||||
|
|
||||||
|
type Cache[V any] interface {
|
||||||
|
Get(ctx context.Context, key string, fetch func(ctx context.Context) (V, error)) (V, error)
|
||||||
|
GetLink(ctx context.Context, key string, fetch func(ctx context.Context) (V, error), link ...string) (V, error)
|
||||||
|
Del(ctx context.Context, key ...string)
|
||||||
|
DelLocal(ctx context.Context, key ...string)
|
||||||
|
Stop()
|
||||||
|
}
|
||||||
|
|
||||||
|
func New[V any](opts ...Option) Cache[V] {
|
||||||
|
opt := defaultOption()
|
||||||
|
for _, o := range opts {
|
||||||
|
o(opt)
|
||||||
|
}
|
||||||
|
|
||||||
|
c := cache[V]{opt: opt}
|
||||||
|
if opt.localSlotNum > 0 && opt.localSlotSize > 0 {
|
||||||
|
createSimpleLRU := func() lru.LRU[string, V] {
|
||||||
|
if opt.expirationEvict {
|
||||||
|
return lru.NewExpirationLRU[string, V](opt.localSlotSize, opt.localSuccessTTL, opt.localFailedTTL, opt.target, c.onEvict)
|
||||||
|
} else {
|
||||||
|
return lru.NewLayLRU[string, V](opt.localSlotSize, opt.localSuccessTTL, opt.localFailedTTL, opt.target, c.onEvict)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if opt.localSlotNum == 1 {
|
||||||
|
c.local = createSimpleLRU()
|
||||||
|
} else {
|
||||||
|
c.local = lru.NewSlotLRU[string, V](opt.localSlotNum, func(key string) uint64 {
|
||||||
|
h := fnv.New64a()
|
||||||
|
h.Write(*(*[]byte)(unsafe.Pointer(&key)))
|
||||||
|
return h.Sum64()
|
||||||
|
}, createSimpleLRU)
|
||||||
|
}
|
||||||
|
if opt.linkSlotNum > 0 {
|
||||||
|
c.link = link.New(opt.linkSlotNum)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return &c
|
||||||
|
}
|
||||||
|
|
||||||
|
type cache[V any] struct {
|
||||||
|
opt *option
|
||||||
|
link link.Link
|
||||||
|
local lru.LRU[string, V]
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *cache[V]) onEvict(key string, value V) {
|
||||||
|
if c.link != nil {
|
||||||
|
lks := c.link.Del(key)
|
||||||
|
for k := range lks {
|
||||||
|
if key != k { // prevent deadlock
|
||||||
|
c.local.Del(k)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *cache[V]) del(key ...string) {
|
||||||
|
if c.local == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
for _, k := range key {
|
||||||
|
c.local.Del(k)
|
||||||
|
if c.link != nil {
|
||||||
|
lks := c.link.Del(k)
|
||||||
|
for k := range lks {
|
||||||
|
c.local.Del(k)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *cache[V]) Get(ctx context.Context, key string, fetch func(ctx context.Context) (V, error)) (V, error) {
|
||||||
|
return c.GetLink(ctx, key, fetch)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *cache[V]) GetLink(ctx context.Context, key string, fetch func(ctx context.Context) (V, error), link ...string) (V, error) {
|
||||||
|
if c.local != nil {
|
||||||
|
return c.local.Get(key, func() (V, error) {
|
||||||
|
if len(link) > 0 {
|
||||||
|
c.link.Link(key, link...)
|
||||||
|
}
|
||||||
|
return fetch(ctx)
|
||||||
|
})
|
||||||
|
} else {
|
||||||
|
return fetch(ctx)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *cache[V]) Del(ctx context.Context, key ...string) {
|
||||||
|
for _, fn := range c.opt.delFn {
|
||||||
|
fn(ctx, key...)
|
||||||
|
}
|
||||||
|
c.del(key...)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *cache[V]) DelLocal(ctx context.Context, key ...string) {
|
||||||
|
c.del(key...)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *cache[V]) Stop() {
|
||||||
|
c.local.Stop()
|
||||||
|
}
|
@ -0,0 +1,79 @@
|
|||||||
|
package localcache
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"math/rand"
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestName(t *testing.T) {
|
||||||
|
c := New[string](WithExpirationEvict())
|
||||||
|
//c := New[string]()
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
const (
|
||||||
|
num = 10000
|
||||||
|
tNum = 10000
|
||||||
|
kNum = 100000
|
||||||
|
pNum = 100
|
||||||
|
)
|
||||||
|
|
||||||
|
getKey := func(v uint64) string {
|
||||||
|
return fmt.Sprintf("key_%d", v%kNum)
|
||||||
|
}
|
||||||
|
|
||||||
|
start := time.Now()
|
||||||
|
t.Log("start", start)
|
||||||
|
|
||||||
|
var (
|
||||||
|
get atomic.Int64
|
||||||
|
del atomic.Int64
|
||||||
|
)
|
||||||
|
|
||||||
|
incrGet := func() {
|
||||||
|
if v := get.Add(1); v%pNum == 0 {
|
||||||
|
//t.Log("#get count", v/pNum)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
incrDel := func() {
|
||||||
|
if v := del.Add(1); v%pNum == 0 {
|
||||||
|
//t.Log("@del count", v/pNum)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
|
||||||
|
for i := 0; i < tNum; i++ {
|
||||||
|
wg.Add(2)
|
||||||
|
go func() {
|
||||||
|
defer wg.Done()
|
||||||
|
for i := 0; i < num; i++ {
|
||||||
|
c.Get(ctx, getKey(rand.Uint64()), func(ctx context.Context) (string, error) {
|
||||||
|
return fmt.Sprintf("index_%d", i), nil
|
||||||
|
})
|
||||||
|
incrGet()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
defer wg.Done()
|
||||||
|
time.Sleep(time.Second / 10)
|
||||||
|
for i := 0; i < num; i++ {
|
||||||
|
c.Del(ctx, getKey(rand.Uint64()))
|
||||||
|
incrDel()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
|
||||||
|
wg.Wait()
|
||||||
|
end := time.Now()
|
||||||
|
t.Log("end", end)
|
||||||
|
t.Log("time", end.Sub(start))
|
||||||
|
t.Log("get", get.Load())
|
||||||
|
t.Log("del", del.Load())
|
||||||
|
// 137.35s
|
||||||
|
}
|
@ -0,0 +1,5 @@
|
|||||||
|
module github.com/openimsdk/localcache
|
||||||
|
|
||||||
|
go 1.19
|
||||||
|
|
||||||
|
require github.com/hashicorp/golang-lru/v2 v2.0.7
|
@ -0,0 +1,109 @@
|
|||||||
|
package link
|
||||||
|
|
||||||
|
import (
|
||||||
|
"hash/fnv"
|
||||||
|
"sync"
|
||||||
|
"unsafe"
|
||||||
|
)
|
||||||
|
|
||||||
|
type Link interface {
|
||||||
|
Link(key string, link ...string)
|
||||||
|
Del(key string) map[string]struct{}
|
||||||
|
}
|
||||||
|
|
||||||
|
func newLinkKey() *linkKey {
|
||||||
|
return &linkKey{
|
||||||
|
data: make(map[string]map[string]struct{}),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type linkKey struct {
|
||||||
|
lock sync.Mutex
|
||||||
|
data map[string]map[string]struct{}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *linkKey) link(key string, link ...string) {
|
||||||
|
x.lock.Lock()
|
||||||
|
defer x.lock.Unlock()
|
||||||
|
v, ok := x.data[key]
|
||||||
|
if !ok {
|
||||||
|
v = make(map[string]struct{})
|
||||||
|
x.data[key] = v
|
||||||
|
}
|
||||||
|
for _, k := range link {
|
||||||
|
v[k] = struct{}{}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *linkKey) del(key string) map[string]struct{} {
|
||||||
|
x.lock.Lock()
|
||||||
|
defer x.lock.Unlock()
|
||||||
|
ks, ok := x.data[key]
|
||||||
|
if !ok {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
delete(x.data, key)
|
||||||
|
return ks
|
||||||
|
}
|
||||||
|
|
||||||
|
func New(n int) Link {
|
||||||
|
if n <= 0 {
|
||||||
|
panic("must be greater than 0")
|
||||||
|
}
|
||||||
|
slots := make([]*linkKey, n)
|
||||||
|
for i := 0; i < len(slots); i++ {
|
||||||
|
slots[i] = newLinkKey()
|
||||||
|
}
|
||||||
|
return &slot{
|
||||||
|
n: uint64(n),
|
||||||
|
slots: slots,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type slot struct {
|
||||||
|
n uint64
|
||||||
|
slots []*linkKey
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *slot) index(s string) uint64 {
|
||||||
|
h := fnv.New64a()
|
||||||
|
_, _ = h.Write(*(*[]byte)(unsafe.Pointer(&s)))
|
||||||
|
return h.Sum64() % x.n
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *slot) Link(key string, link ...string) {
|
||||||
|
if len(link) == 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
mk := key
|
||||||
|
lks := make([]string, len(link))
|
||||||
|
for i, k := range link {
|
||||||
|
lks[i] = k
|
||||||
|
}
|
||||||
|
x.slots[x.index(mk)].link(mk, lks...)
|
||||||
|
for _, lk := range lks {
|
||||||
|
x.slots[x.index(lk)].link(lk, mk)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *slot) Del(key string) map[string]struct{} {
|
||||||
|
return x.delKey(key)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *slot) delKey(k string) map[string]struct{} {
|
||||||
|
del := make(map[string]struct{})
|
||||||
|
stack := []string{k}
|
||||||
|
for len(stack) > 0 {
|
||||||
|
curr := stack[len(stack)-1]
|
||||||
|
stack = stack[:len(stack)-1]
|
||||||
|
if _, ok := del[curr]; ok {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
del[curr] = struct{}{}
|
||||||
|
childKeys := x.slots[x.index(curr)].del(curr)
|
||||||
|
for ck := range childKeys {
|
||||||
|
stack = append(stack, ck)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return del
|
||||||
|
}
|
@ -0,0 +1,20 @@
|
|||||||
|
package link
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestName(t *testing.T) {
|
||||||
|
|
||||||
|
v := New(1)
|
||||||
|
|
||||||
|
//v.Link("a:1", "b:1", "c:1", "d:1")
|
||||||
|
v.Link("a:1", "b:1", "c:1")
|
||||||
|
v.Link("z:1", "b:1")
|
||||||
|
|
||||||
|
//v.DelKey("a:1")
|
||||||
|
v.Del("z:1")
|
||||||
|
|
||||||
|
t.Log(v)
|
||||||
|
|
||||||
|
}
|
@ -0,0 +1,20 @@
|
|||||||
|
package lru
|
||||||
|
|
||||||
|
import "github.com/hashicorp/golang-lru/v2/simplelru"
|
||||||
|
|
||||||
|
type EvictCallback[K comparable, V any] simplelru.EvictCallback[K, V]
|
||||||
|
|
||||||
|
type LRU[K comparable, V any] interface {
|
||||||
|
Get(key K, fetch func() (V, error)) (V, error)
|
||||||
|
Del(key K) bool
|
||||||
|
Stop()
|
||||||
|
}
|
||||||
|
|
||||||
|
type Target interface {
|
||||||
|
IncrGetHit()
|
||||||
|
IncrGetSuccess()
|
||||||
|
IncrGetFailed()
|
||||||
|
|
||||||
|
IncrDelHit()
|
||||||
|
IncrDelNotFound()
|
||||||
|
}
|
@ -0,0 +1,78 @@
|
|||||||
|
package lru
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/hashicorp/golang-lru/v2/expirable"
|
||||||
|
"sync"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
func NewExpirationLRU[K comparable, V any](size int, successTTL, failedTTL time.Duration, target Target, onEvict EvictCallback[K, V]) LRU[K, V] {
|
||||||
|
var cb expirable.EvictCallback[K, *expirationLruItem[V]]
|
||||||
|
if onEvict != nil {
|
||||||
|
cb = func(key K, value *expirationLruItem[V]) {
|
||||||
|
onEvict(key, value.value)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
core := expirable.NewLRU[K, *expirationLruItem[V]](size, cb, successTTL)
|
||||||
|
return &ExpirationLRU[K, V]{
|
||||||
|
core: core,
|
||||||
|
successTTL: successTTL,
|
||||||
|
failedTTL: failedTTL,
|
||||||
|
target: target,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type expirationLruItem[V any] struct {
|
||||||
|
lock sync.RWMutex
|
||||||
|
err error
|
||||||
|
value V
|
||||||
|
}
|
||||||
|
|
||||||
|
type ExpirationLRU[K comparable, V any] struct {
|
||||||
|
lock sync.Mutex
|
||||||
|
core *expirable.LRU[K, *expirationLruItem[V]]
|
||||||
|
successTTL time.Duration
|
||||||
|
failedTTL time.Duration
|
||||||
|
target Target
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *ExpirationLRU[K, V]) Get(key K, fetch func() (V, error)) (V, error) {
|
||||||
|
x.lock.Lock()
|
||||||
|
v, ok := x.core.Get(key)
|
||||||
|
if ok {
|
||||||
|
x.lock.Unlock()
|
||||||
|
x.target.IncrGetSuccess()
|
||||||
|
v.lock.RLock()
|
||||||
|
defer v.lock.RUnlock()
|
||||||
|
return v.value, v.err
|
||||||
|
} else {
|
||||||
|
v = &expirationLruItem[V]{}
|
||||||
|
x.core.Add(key, v)
|
||||||
|
v.lock.Lock()
|
||||||
|
x.lock.Unlock()
|
||||||
|
defer v.lock.Unlock()
|
||||||
|
v.value, v.err = fetch()
|
||||||
|
if v.err == nil {
|
||||||
|
x.target.IncrGetSuccess()
|
||||||
|
} else {
|
||||||
|
x.target.IncrGetFailed()
|
||||||
|
x.core.Remove(key)
|
||||||
|
}
|
||||||
|
return v.value, v.err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *ExpirationLRU[K, V]) Del(key K) bool {
|
||||||
|
x.lock.Lock()
|
||||||
|
ok := x.core.Remove(key)
|
||||||
|
x.lock.Unlock()
|
||||||
|
if ok {
|
||||||
|
x.target.IncrDelHit()
|
||||||
|
} else {
|
||||||
|
x.target.IncrDelNotFound()
|
||||||
|
}
|
||||||
|
return ok
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *ExpirationLRU[K, V]) Stop() {
|
||||||
|
}
|
@ -0,0 +1,90 @@
|
|||||||
|
package lru
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/hashicorp/golang-lru/v2/simplelru"
|
||||||
|
"sync"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
type layLruItem[V any] struct {
|
||||||
|
lock sync.Mutex
|
||||||
|
expires int64
|
||||||
|
err error
|
||||||
|
value V
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewLayLRU[K comparable, V any](size int, successTTL, failedTTL time.Duration, target Target, onEvict EvictCallback[K, V]) *LayLRU[K, V] {
|
||||||
|
var cb simplelru.EvictCallback[K, *layLruItem[V]]
|
||||||
|
if onEvict != nil {
|
||||||
|
cb = func(key K, value *layLruItem[V]) {
|
||||||
|
onEvict(key, value.value)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
core, err := simplelru.NewLRU[K, *layLruItem[V]](size, cb)
|
||||||
|
if err != nil {
|
||||||
|
panic(err)
|
||||||
|
}
|
||||||
|
return &LayLRU[K, V]{
|
||||||
|
core: core,
|
||||||
|
successTTL: successTTL,
|
||||||
|
failedTTL: failedTTL,
|
||||||
|
target: target,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type LayLRU[K comparable, V any] struct {
|
||||||
|
lock sync.Mutex
|
||||||
|
core *simplelru.LRU[K, *layLruItem[V]]
|
||||||
|
successTTL time.Duration
|
||||||
|
failedTTL time.Duration
|
||||||
|
target Target
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *LayLRU[K, V]) Get(key K, fetch func() (V, error)) (V, error) {
|
||||||
|
x.lock.Lock()
|
||||||
|
v, ok := x.core.Get(key)
|
||||||
|
if ok {
|
||||||
|
x.lock.Unlock()
|
||||||
|
v.lock.Lock()
|
||||||
|
expires, value, err := v.expires, v.value, v.err
|
||||||
|
if expires != 0 && expires > time.Now().UnixMilli() {
|
||||||
|
v.lock.Unlock()
|
||||||
|
x.target.IncrGetHit()
|
||||||
|
return value, err
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
v = &layLruItem[V]{}
|
||||||
|
x.core.Add(key, v)
|
||||||
|
v.lock.Lock()
|
||||||
|
x.lock.Unlock()
|
||||||
|
}
|
||||||
|
defer v.lock.Unlock()
|
||||||
|
if v.expires > time.Now().UnixMilli() {
|
||||||
|
return v.value, v.err
|
||||||
|
}
|
||||||
|
v.value, v.err = fetch()
|
||||||
|
if v.err == nil {
|
||||||
|
v.expires = time.Now().Add(x.successTTL).UnixMilli()
|
||||||
|
x.target.IncrGetSuccess()
|
||||||
|
} else {
|
||||||
|
v.expires = time.Now().Add(x.failedTTL).UnixMilli()
|
||||||
|
x.target.IncrGetFailed()
|
||||||
|
}
|
||||||
|
return v.value, v.err
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *LayLRU[K, V]) Del(key K) bool {
|
||||||
|
x.lock.Lock()
|
||||||
|
ok := x.core.Remove(key)
|
||||||
|
x.lock.Unlock()
|
||||||
|
if ok {
|
||||||
|
x.target.IncrDelHit()
|
||||||
|
} else {
|
||||||
|
x.target.IncrDelNotFound()
|
||||||
|
}
|
||||||
|
return ok
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *LayLRU[K, V]) Stop() {
|
||||||
|
|
||||||
|
}
|
@ -0,0 +1,104 @@
|
|||||||
|
package lru
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"hash/fnv"
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
"unsafe"
|
||||||
|
)
|
||||||
|
|
||||||
|
type cacheTarget struct {
|
||||||
|
getHit int64
|
||||||
|
getSuccess int64
|
||||||
|
getFailed int64
|
||||||
|
delHit int64
|
||||||
|
delNotFound int64
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *cacheTarget) IncrGetHit() {
|
||||||
|
atomic.AddInt64(&r.getHit, 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *cacheTarget) IncrGetSuccess() {
|
||||||
|
atomic.AddInt64(&r.getSuccess, 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *cacheTarget) IncrGetFailed() {
|
||||||
|
atomic.AddInt64(&r.getFailed, 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *cacheTarget) IncrDelHit() {
|
||||||
|
atomic.AddInt64(&r.delHit, 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *cacheTarget) IncrDelNotFound() {
|
||||||
|
atomic.AddInt64(&r.delNotFound, 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *cacheTarget) String() string {
|
||||||
|
return fmt.Sprintf("getHit: %d, getSuccess: %d, getFailed: %d, delHit: %d, delNotFound: %d", r.getHit, r.getSuccess, r.getFailed, r.delHit, r.delNotFound)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestName(t *testing.T) {
|
||||||
|
target := &cacheTarget{}
|
||||||
|
l := NewSlotLRU[string, string](100, func(k string) uint64 {
|
||||||
|
h := fnv.New64a()
|
||||||
|
h.Write(*(*[]byte)(unsafe.Pointer(&k)))
|
||||||
|
return h.Sum64()
|
||||||
|
}, func() LRU[string, string] {
|
||||||
|
return NewExpirationLRU[string, string](100, time.Second*60, time.Second, target, nil)
|
||||||
|
})
|
||||||
|
//l := NewInertiaLRU[string, string](1000, time.Second*20, time.Second*5, target)
|
||||||
|
|
||||||
|
fn := func(key string, n int, fetch func() (string, error)) {
|
||||||
|
for i := 0; i < n; i++ {
|
||||||
|
//v, err := l.Get(key, fetch)
|
||||||
|
//if err == nil {
|
||||||
|
// t.Log("key", key, "value", v)
|
||||||
|
//} else {
|
||||||
|
// t.Error("key", key, err)
|
||||||
|
//}
|
||||||
|
v, err := l.Get(key, fetch)
|
||||||
|
//time.Sleep(time.Second / 100)
|
||||||
|
func(v ...any) {}(v, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
tmp := make(map[string]struct{})
|
||||||
|
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
for i := 0; i < 10000; i++ {
|
||||||
|
wg.Add(1)
|
||||||
|
key := fmt.Sprintf("key_%d", i%200)
|
||||||
|
tmp[key] = struct{}{}
|
||||||
|
go func() {
|
||||||
|
defer wg.Done()
|
||||||
|
//t.Log(key)
|
||||||
|
fn(key, 10000, func() (string, error) {
|
||||||
|
//time.Sleep(time.Second * 3)
|
||||||
|
//t.Log(time.Now(), "key", key, "fetch")
|
||||||
|
//if rand.Uint32()%5 == 0 {
|
||||||
|
// return "value_" + key, nil
|
||||||
|
//}
|
||||||
|
//return "", errors.New("rand error")
|
||||||
|
return "value_" + key, nil
|
||||||
|
})
|
||||||
|
}()
|
||||||
|
|
||||||
|
//wg.Add(1)
|
||||||
|
//go func() {
|
||||||
|
// defer wg.Done()
|
||||||
|
// for i := 0; i < 10; i++ {
|
||||||
|
// l.Del(key)
|
||||||
|
// time.Sleep(time.Second / 3)
|
||||||
|
// }
|
||||||
|
//}()
|
||||||
|
}
|
||||||
|
wg.Wait()
|
||||||
|
t.Log(len(tmp))
|
||||||
|
t.Log(target.String())
|
||||||
|
|
||||||
|
}
|
@ -0,0 +1,37 @@
|
|||||||
|
package lru
|
||||||
|
|
||||||
|
func NewSlotLRU[K comparable, V any](slotNum int, hash func(K) uint64, create func() LRU[K, V]) LRU[K, V] {
|
||||||
|
x := &slotLRU[K, V]{
|
||||||
|
n: uint64(slotNum),
|
||||||
|
slots: make([]LRU[K, V], slotNum),
|
||||||
|
hash: hash,
|
||||||
|
}
|
||||||
|
for i := 0; i < slotNum; i++ {
|
||||||
|
x.slots[i] = create()
|
||||||
|
}
|
||||||
|
return x
|
||||||
|
}
|
||||||
|
|
||||||
|
type slotLRU[K comparable, V any] struct {
|
||||||
|
n uint64
|
||||||
|
slots []LRU[K, V]
|
||||||
|
hash func(k K) uint64
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *slotLRU[K, V]) getIndex(k K) uint64 {
|
||||||
|
return x.hash(k) % x.n
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *slotLRU[K, V]) Get(key K, fetch func() (V, error)) (V, error) {
|
||||||
|
return x.slots[x.getIndex(key)].Get(key, fetch)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *slotLRU[K, V]) Del(key K) bool {
|
||||||
|
return x.slots[x.getIndex(key)].Del(key)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (x *slotLRU[K, V]) Stop() {
|
||||||
|
for _, slot := range x.slots {
|
||||||
|
slot.Stop()
|
||||||
|
}
|
||||||
|
}
|
@ -0,0 +1,121 @@
|
|||||||
|
package localcache
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"github.com/openimsdk/localcache/lru"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
func defaultOption() *option {
|
||||||
|
return &option{
|
||||||
|
localSlotNum: 500,
|
||||||
|
localSlotSize: 20000,
|
||||||
|
linkSlotNum: 500,
|
||||||
|
expirationEvict: false,
|
||||||
|
localSuccessTTL: time.Minute,
|
||||||
|
localFailedTTL: time.Second * 5,
|
||||||
|
delFn: make([]func(ctx context.Context, key ...string), 0, 2),
|
||||||
|
target: emptyTarget{},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type option struct {
|
||||||
|
localSlotNum int
|
||||||
|
localSlotSize int
|
||||||
|
linkSlotNum int
|
||||||
|
// expirationEvict: true means that the cache will be actively cleared when the timer expires,
|
||||||
|
// false means that the cache will be lazily deleted.
|
||||||
|
expirationEvict bool
|
||||||
|
localSuccessTTL time.Duration
|
||||||
|
localFailedTTL time.Duration
|
||||||
|
delFn []func(ctx context.Context, key ...string)
|
||||||
|
target lru.Target
|
||||||
|
}
|
||||||
|
|
||||||
|
type Option func(o *option)
|
||||||
|
|
||||||
|
func WithExpirationEvict() Option {
|
||||||
|
return func(o *option) {
|
||||||
|
o.expirationEvict = true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func WithLazy() Option {
|
||||||
|
return func(o *option) {
|
||||||
|
o.expirationEvict = false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func WithLocalDisable() Option {
|
||||||
|
return WithLinkSlotNum(0)
|
||||||
|
}
|
||||||
|
|
||||||
|
func WithLinkDisable() Option {
|
||||||
|
return WithLinkSlotNum(0)
|
||||||
|
}
|
||||||
|
|
||||||
|
func WithLinkSlotNum(linkSlotNum int) Option {
|
||||||
|
return func(o *option) {
|
||||||
|
o.linkSlotNum = linkSlotNum
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func WithLocalSlotNum(localSlotNum int) Option {
|
||||||
|
return func(o *option) {
|
||||||
|
o.localSlotNum = localSlotNum
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func WithLocalSlotSize(localSlotSize int) Option {
|
||||||
|
return func(o *option) {
|
||||||
|
o.localSlotSize = localSlotSize
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func WithLocalSuccessTTL(localSuccessTTL time.Duration) Option {
|
||||||
|
if localSuccessTTL < 0 {
|
||||||
|
panic("localSuccessTTL should be greater than 0")
|
||||||
|
}
|
||||||
|
return func(o *option) {
|
||||||
|
o.localSuccessTTL = localSuccessTTL
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func WithLocalFailedTTL(localFailedTTL time.Duration) Option {
|
||||||
|
if localFailedTTL < 0 {
|
||||||
|
panic("localFailedTTL should be greater than 0")
|
||||||
|
}
|
||||||
|
return func(o *option) {
|
||||||
|
o.localFailedTTL = localFailedTTL
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func WithTarget(target lru.Target) Option {
|
||||||
|
if target == nil {
|
||||||
|
panic("target should not be nil")
|
||||||
|
}
|
||||||
|
return func(o *option) {
|
||||||
|
o.target = target
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func WithDeleteKeyBefore(fn func(ctx context.Context, key ...string)) Option {
|
||||||
|
if fn == nil {
|
||||||
|
panic("fn should not be nil")
|
||||||
|
}
|
||||||
|
return func(o *option) {
|
||||||
|
o.delFn = append(o.delFn, fn)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type emptyTarget struct{}
|
||||||
|
|
||||||
|
func (e emptyTarget) IncrGetHit() {}
|
||||||
|
|
||||||
|
func (e emptyTarget) IncrGetSuccess() {}
|
||||||
|
|
||||||
|
func (e emptyTarget) IncrGetFailed() {}
|
||||||
|
|
||||||
|
func (e emptyTarget) IncrDelHit() {}
|
||||||
|
|
||||||
|
func (e emptyTarget) IncrDelNotFound() {}
|
@ -0,0 +1,9 @@
|
|||||||
|
package localcache
|
||||||
|
|
||||||
|
func AnyValue[V any](v any, err error) (V, error) {
|
||||||
|
if err != nil {
|
||||||
|
var zero V
|
||||||
|
return zero, err
|
||||||
|
}
|
||||||
|
return v.(V), nil
|
||||||
|
}
|
@ -0,0 +1,20 @@
|
|||||||
|
package rpccache
|
||||||
|
|
||||||
|
func newListMap[V comparable](values []V, err error) (*listMap[V], error) {
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
lm := &listMap[V]{
|
||||||
|
List: values,
|
||||||
|
Map: make(map[V]struct{}, len(values)),
|
||||||
|
}
|
||||||
|
for _, value := range values {
|
||||||
|
lm.Map[value] = struct{}{}
|
||||||
|
}
|
||||||
|
return lm, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
type listMap[V comparable] struct {
|
||||||
|
List []V
|
||||||
|
Map map[V]struct{}
|
||||||
|
}
|
@ -0,0 +1,112 @@
|
|||||||
|
package rpccache
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
pbconversation "github.com/OpenIMSDK/protocol/conversation"
|
||||||
|
"github.com/OpenIMSDK/tools/errs"
|
||||||
|
"github.com/OpenIMSDK/tools/log"
|
||||||
|
"github.com/openimsdk/localcache"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/cachekey"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/config"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/rpcclient"
|
||||||
|
"github.com/redis/go-redis/v9"
|
||||||
|
)
|
||||||
|
|
||||||
|
func NewConversationLocalCache(client rpcclient.ConversationRpcClient, cli redis.UniversalClient) *ConversationLocalCache {
|
||||||
|
lc := config.Config.LocalCache.Conversation
|
||||||
|
log.ZDebug(context.Background(), "ConversationLocalCache", "topic", lc.Topic, "slotNum", lc.SlotNum, "slotSize", lc.SlotSize, "enable", lc.Enable())
|
||||||
|
x := &ConversationLocalCache{
|
||||||
|
client: client,
|
||||||
|
local: localcache.New[any](
|
||||||
|
localcache.WithLocalSlotNum(lc.SlotNum),
|
||||||
|
localcache.WithLocalSlotSize(lc.SlotSize),
|
||||||
|
localcache.WithLinkSlotNum(lc.SlotNum),
|
||||||
|
localcache.WithLocalSuccessTTL(lc.Success()),
|
||||||
|
localcache.WithLocalFailedTTL(lc.Failed()),
|
||||||
|
),
|
||||||
|
}
|
||||||
|
if lc.Enable() {
|
||||||
|
go subscriberRedisDeleteCache(context.Background(), cli, lc.Topic, x.local.DelLocal)
|
||||||
|
}
|
||||||
|
return x
|
||||||
|
}
|
||||||
|
|
||||||
|
type ConversationLocalCache struct {
|
||||||
|
client rpcclient.ConversationRpcClient
|
||||||
|
local localcache.Cache[any]
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *ConversationLocalCache) GetConversationIDs(ctx context.Context, ownerUserID string) (val []string, err error) {
|
||||||
|
log.ZDebug(ctx, "ConversationLocalCache GetConversationIDs req", "ownerUserID", ownerUserID)
|
||||||
|
defer func() {
|
||||||
|
if err == nil {
|
||||||
|
log.ZDebug(ctx, "ConversationLocalCache GetConversationIDs return", "value", val)
|
||||||
|
} else {
|
||||||
|
log.ZError(ctx, "ConversationLocalCache GetConversationIDs return", err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
return localcache.AnyValue[[]string](c.local.Get(ctx, cachekey.GetConversationIDsKey(ownerUserID), func(ctx context.Context) (any, error) {
|
||||||
|
log.ZDebug(ctx, "ConversationLocalCache GetConversationIDs rpc", "ownerUserID", ownerUserID)
|
||||||
|
return c.client.GetConversationIDs(ctx, ownerUserID)
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *ConversationLocalCache) GetConversation(ctx context.Context, userID, conversationID string) (val *pbconversation.Conversation, err error) {
|
||||||
|
log.ZDebug(ctx, "ConversationLocalCache GetConversation req", "userID", userID, "conversationID", conversationID)
|
||||||
|
defer func() {
|
||||||
|
if err == nil {
|
||||||
|
log.ZDebug(ctx, "ConversationLocalCache GetConversation return", "value", val)
|
||||||
|
} else {
|
||||||
|
log.ZError(ctx, "ConversationLocalCache GetConversation return", err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
return localcache.AnyValue[*pbconversation.Conversation](c.local.Get(ctx, cachekey.GetConversationKey(userID, conversationID), func(ctx context.Context) (any, error) {
|
||||||
|
log.ZDebug(ctx, "ConversationLocalCache GetConversation rpc", "userID", userID, "conversationID", conversationID)
|
||||||
|
return c.client.GetConversation(ctx, userID, conversationID)
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *ConversationLocalCache) GetSingleConversationRecvMsgOpt(ctx context.Context, userID, conversationID string) (int32, error) {
|
||||||
|
conv, err := c.GetConversation(ctx, userID, conversationID)
|
||||||
|
if err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
return conv.RecvMsgOpt, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *ConversationLocalCache) GetConversations(ctx context.Context, ownerUserID string, conversationIDs []string) ([]*pbconversation.Conversation, error) {
|
||||||
|
conversations := make([]*pbconversation.Conversation, 0, len(conversationIDs))
|
||||||
|
for _, conversationID := range conversationIDs {
|
||||||
|
conversation, err := c.GetConversation(ctx, ownerUserID, conversationID)
|
||||||
|
if err != nil {
|
||||||
|
if errs.ErrRecordNotFound.Is(err) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
conversations = append(conversations, conversation)
|
||||||
|
}
|
||||||
|
return conversations, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *ConversationLocalCache) getConversationNotReceiveMessageUserIDs(ctx context.Context, conversationID string) (*listMap[string], error) {
|
||||||
|
return localcache.AnyValue[*listMap[string]](c.local.Get(ctx, cachekey.GetConversationNotReceiveMessageUserIDsKey(conversationID), func(ctx context.Context) (any, error) {
|
||||||
|
return newListMap(c.client.GetConversationNotReceiveMessageUserIDs(ctx, conversationID))
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *ConversationLocalCache) GetConversationNotReceiveMessageUserIDs(ctx context.Context, conversationID string) ([]string, error) {
|
||||||
|
res, err := c.getConversationNotReceiveMessageUserIDs(ctx, conversationID)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return res.List, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *ConversationLocalCache) GetConversationNotReceiveMessageUserIDMap(ctx context.Context, conversationID string) (map[string]struct{}, error) {
|
||||||
|
res, err := c.getConversationNotReceiveMessageUserIDs(ctx, conversationID)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return res.Map, nil
|
||||||
|
}
|
@ -0,0 +1,66 @@
|
|||||||
|
package rpccache
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"github.com/OpenIMSDK/tools/log"
|
||||||
|
"github.com/openimsdk/localcache"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/cachekey"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/config"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/rpcclient"
|
||||||
|
"github.com/redis/go-redis/v9"
|
||||||
|
)
|
||||||
|
|
||||||
|
func NewFriendLocalCache(client rpcclient.FriendRpcClient, cli redis.UniversalClient) *FriendLocalCache {
|
||||||
|
lc := config.Config.LocalCache.Friend
|
||||||
|
log.ZDebug(context.Background(), "FriendLocalCache", "topic", lc.Topic, "slotNum", lc.SlotNum, "slotSize", lc.SlotSize, "enable", lc.Enable())
|
||||||
|
x := &FriendLocalCache{
|
||||||
|
client: client,
|
||||||
|
local: localcache.New[any](
|
||||||
|
localcache.WithLocalSlotNum(lc.SlotNum),
|
||||||
|
localcache.WithLocalSlotSize(lc.SlotSize),
|
||||||
|
localcache.WithLinkSlotNum(lc.SlotNum),
|
||||||
|
localcache.WithLocalSuccessTTL(lc.Success()),
|
||||||
|
localcache.WithLocalFailedTTL(lc.Failed()),
|
||||||
|
),
|
||||||
|
}
|
||||||
|
if lc.Enable() {
|
||||||
|
go subscriberRedisDeleteCache(context.Background(), cli, lc.Topic, x.local.DelLocal)
|
||||||
|
}
|
||||||
|
return x
|
||||||
|
}
|
||||||
|
|
||||||
|
type FriendLocalCache struct {
|
||||||
|
client rpcclient.FriendRpcClient
|
||||||
|
local localcache.Cache[any]
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *FriendLocalCache) IsFriend(ctx context.Context, possibleFriendUserID, userID string) (val bool, err error) {
|
||||||
|
log.ZDebug(ctx, "FriendLocalCache IsFriend req", "possibleFriendUserID", possibleFriendUserID, "userID", userID)
|
||||||
|
defer func() {
|
||||||
|
if err == nil {
|
||||||
|
log.ZDebug(ctx, "FriendLocalCache IsFriend return", "value", val)
|
||||||
|
} else {
|
||||||
|
log.ZError(ctx, "FriendLocalCache IsFriend return", err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
return localcache.AnyValue[bool](f.local.GetLink(ctx, cachekey.GetIsFriendKey(possibleFriendUserID, userID), func(ctx context.Context) (any, error) {
|
||||||
|
log.ZDebug(ctx, "FriendLocalCache IsFriend rpc", "possibleFriendUserID", possibleFriendUserID, "userID", userID)
|
||||||
|
return f.client.IsFriend(ctx, possibleFriendUserID, userID)
|
||||||
|
}, cachekey.GetFriendIDsKey(possibleFriendUserID)))
|
||||||
|
}
|
||||||
|
|
||||||
|
// IsBlack possibleBlackUserID selfUserID
|
||||||
|
func (f *FriendLocalCache) IsBlack(ctx context.Context, possibleBlackUserID, userID string) (val bool, err error) {
|
||||||
|
log.ZDebug(ctx, "FriendLocalCache IsBlack req", "possibleBlackUserID", possibleBlackUserID, "userID", userID)
|
||||||
|
defer func() {
|
||||||
|
if err == nil {
|
||||||
|
log.ZDebug(ctx, "FriendLocalCache IsBlack return", "value", val)
|
||||||
|
} else {
|
||||||
|
log.ZError(ctx, "FriendLocalCache IsBlack return", err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
return localcache.AnyValue[bool](f.local.GetLink(ctx, cachekey.GetIsBlackIDsKey(possibleBlackUserID, userID), func(ctx context.Context) (any, error) {
|
||||||
|
log.ZDebug(ctx, "FriendLocalCache IsBlack rpc", "possibleBlackUserID", possibleBlackUserID, "userID", userID)
|
||||||
|
return f.client.IsBlack(ctx, possibleBlackUserID, userID)
|
||||||
|
}, cachekey.GetBlackIDsKey(userID)))
|
||||||
|
}
|
@ -0,0 +1,143 @@
|
|||||||
|
package rpccache
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"github.com/OpenIMSDK/protocol/sdkws"
|
||||||
|
"github.com/OpenIMSDK/tools/errs"
|
||||||
|
"github.com/OpenIMSDK/tools/log"
|
||||||
|
"github.com/openimsdk/localcache"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/cachekey"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/config"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/rpcclient"
|
||||||
|
"github.com/redis/go-redis/v9"
|
||||||
|
)
|
||||||
|
|
||||||
|
func NewGroupLocalCache(client rpcclient.GroupRpcClient, cli redis.UniversalClient) *GroupLocalCache {
|
||||||
|
lc := config.Config.LocalCache.Group
|
||||||
|
log.ZDebug(context.Background(), "GroupLocalCache", "topic", lc.Topic, "slotNum", lc.SlotNum, "slotSize", lc.SlotSize, "enable", lc.Enable())
|
||||||
|
x := &GroupLocalCache{
|
||||||
|
client: client,
|
||||||
|
local: localcache.New[any](
|
||||||
|
localcache.WithLocalSlotNum(lc.SlotNum),
|
||||||
|
localcache.WithLocalSlotSize(lc.SlotSize),
|
||||||
|
localcache.WithLinkSlotNum(lc.SlotNum),
|
||||||
|
localcache.WithLocalSuccessTTL(lc.Success()),
|
||||||
|
localcache.WithLocalFailedTTL(lc.Failed()),
|
||||||
|
),
|
||||||
|
}
|
||||||
|
if lc.Enable() {
|
||||||
|
go subscriberRedisDeleteCache(context.Background(), cli, lc.Topic, x.local.DelLocal)
|
||||||
|
}
|
||||||
|
return x
|
||||||
|
}
|
||||||
|
|
||||||
|
type GroupLocalCache struct {
|
||||||
|
client rpcclient.GroupRpcClient
|
||||||
|
local localcache.Cache[any]
|
||||||
|
}
|
||||||
|
|
||||||
|
func (g *GroupLocalCache) getGroupMemberIDs(ctx context.Context, groupID string) (val *listMap[string], err error) {
|
||||||
|
log.ZDebug(ctx, "GroupLocalCache getGroupMemberIDs req", "groupID", groupID)
|
||||||
|
defer func() {
|
||||||
|
if err == nil {
|
||||||
|
log.ZDebug(ctx, "GroupLocalCache getGroupMemberIDs return", "value", val)
|
||||||
|
} else {
|
||||||
|
log.ZError(ctx, "GroupLocalCache getGroupMemberIDs return", err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
return localcache.AnyValue[*listMap[string]](g.local.Get(ctx, cachekey.GetGroupMemberIDsKey(groupID), func(ctx context.Context) (any, error) {
|
||||||
|
log.ZDebug(ctx, "GroupLocalCache getGroupMemberIDs rpc", "groupID", groupID)
|
||||||
|
return newListMap(g.client.GetGroupMemberIDs(ctx, groupID))
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
|
func (g *GroupLocalCache) GetGroupMember(ctx context.Context, groupID, userID string) (val *sdkws.GroupMemberFullInfo, err error) {
|
||||||
|
log.ZDebug(ctx, "GroupLocalCache GetGroupInfo req", "groupID", groupID, "userID", userID)
|
||||||
|
defer func() {
|
||||||
|
if err == nil {
|
||||||
|
log.ZDebug(ctx, "GroupLocalCache GetGroupInfo return", "value", val)
|
||||||
|
} else {
|
||||||
|
log.ZError(ctx, "GroupLocalCache GetGroupInfo return", err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
return localcache.AnyValue[*sdkws.GroupMemberFullInfo](g.local.Get(ctx, cachekey.GetGroupMemberInfoKey(groupID, userID), func(ctx context.Context) (any, error) {
|
||||||
|
log.ZDebug(ctx, "GroupLocalCache GetGroupInfo rpc", "groupID", groupID, "userID", userID)
|
||||||
|
return g.client.GetGroupMemberCache(ctx, groupID, userID)
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
|
func (g *GroupLocalCache) GetGroupInfo(ctx context.Context, groupID string) (val *sdkws.GroupInfo, err error) {
|
||||||
|
log.ZDebug(ctx, "GroupLocalCache GetGroupInfo req", "groupID", groupID)
|
||||||
|
defer func() {
|
||||||
|
if err == nil {
|
||||||
|
log.ZDebug(ctx, "GroupLocalCache GetGroupInfo return", "value", val)
|
||||||
|
} else {
|
||||||
|
log.ZError(ctx, "GroupLocalCache GetGroupInfo return", err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
return localcache.AnyValue[*sdkws.GroupInfo](g.local.Get(ctx, cachekey.GetGroupInfoKey(groupID), func(ctx context.Context) (any, error) {
|
||||||
|
log.ZDebug(ctx, "GroupLocalCache GetGroupInfo rpc", "groupID", groupID)
|
||||||
|
return g.client.GetGroupInfoCache(ctx, groupID)
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
|
func (g *GroupLocalCache) GetGroupMemberIDs(ctx context.Context, groupID string) ([]string, error) {
|
||||||
|
res, err := g.getGroupMemberIDs(ctx, groupID)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return res.List, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (g *GroupLocalCache) GetGroupMemberIDMap(ctx context.Context, groupID string) (map[string]struct{}, error) {
|
||||||
|
res, err := g.getGroupMemberIDs(ctx, groupID)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return res.Map, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (g *GroupLocalCache) GetGroupInfos(ctx context.Context, groupIDs []string) ([]*sdkws.GroupInfo, error) {
|
||||||
|
groupInfos := make([]*sdkws.GroupInfo, 0, len(groupIDs))
|
||||||
|
for _, groupID := range groupIDs {
|
||||||
|
groupInfo, err := g.GetGroupInfo(ctx, groupID)
|
||||||
|
if err != nil {
|
||||||
|
if errs.ErrRecordNotFound.Is(err) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
groupInfos = append(groupInfos, groupInfo)
|
||||||
|
}
|
||||||
|
return groupInfos, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (g *GroupLocalCache) GetGroupMembers(ctx context.Context, groupID string, userIDs []string) ([]*sdkws.GroupMemberFullInfo, error) {
|
||||||
|
members := make([]*sdkws.GroupMemberFullInfo, 0, len(userIDs))
|
||||||
|
for _, userID := range userIDs {
|
||||||
|
member, err := g.GetGroupMember(ctx, groupID, userID)
|
||||||
|
if err != nil {
|
||||||
|
if errs.ErrRecordNotFound.Is(err) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
members = append(members, member)
|
||||||
|
}
|
||||||
|
return members, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (g *GroupLocalCache) GetGroupMemberInfoMap(ctx context.Context, groupID string, userIDs []string) (map[string]*sdkws.GroupMemberFullInfo, error) {
|
||||||
|
members := make(map[string]*sdkws.GroupMemberFullInfo)
|
||||||
|
for _, userID := range userIDs {
|
||||||
|
member, err := g.GetGroupMember(ctx, groupID, userID)
|
||||||
|
if err != nil {
|
||||||
|
if errs.ErrRecordNotFound.Is(err) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
members[userID] = member
|
||||||
|
}
|
||||||
|
return members, nil
|
||||||
|
}
|
@ -0,0 +1,23 @@
|
|||||||
|
package rpccache
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"github.com/OpenIMSDK/tools/log"
|
||||||
|
"github.com/redis/go-redis/v9"
|
||||||
|
)
|
||||||
|
|
||||||
|
func subscriberRedisDeleteCache(ctx context.Context, client redis.UniversalClient, channel string, del func(ctx context.Context, key ...string)) {
|
||||||
|
for message := range client.Subscribe(ctx, channel).Channel() {
|
||||||
|
log.ZDebug(ctx, "subscriberRedisDeleteCache", "channel", channel, "payload", message.Payload)
|
||||||
|
var keys []string
|
||||||
|
if err := json.Unmarshal([]byte(message.Payload), &keys); err != nil {
|
||||||
|
log.ZError(ctx, "subscriberRedisDeleteCache json.Unmarshal error", err)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if len(keys) == 0 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
del(ctx, keys...)
|
||||||
|
}
|
||||||
|
}
|
@ -0,0 +1,97 @@
|
|||||||
|
package rpccache
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"github.com/OpenIMSDK/protocol/sdkws"
|
||||||
|
"github.com/OpenIMSDK/tools/errs"
|
||||||
|
"github.com/OpenIMSDK/tools/log"
|
||||||
|
"github.com/openimsdk/localcache"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/cachekey"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/common/config"
|
||||||
|
"github.com/openimsdk/open-im-server/v3/pkg/rpcclient"
|
||||||
|
"github.com/redis/go-redis/v9"
|
||||||
|
)
|
||||||
|
|
||||||
|
func NewUserLocalCache(client rpcclient.UserRpcClient, cli redis.UniversalClient) *UserLocalCache {
|
||||||
|
lc := config.Config.LocalCache.User
|
||||||
|
log.ZDebug(context.Background(), "UserLocalCache", "topic", lc.Topic, "slotNum", lc.SlotNum, "slotSize", lc.SlotSize, "enable", lc.Enable())
|
||||||
|
x := &UserLocalCache{
|
||||||
|
client: client,
|
||||||
|
local: localcache.New[any](
|
||||||
|
localcache.WithLocalSlotNum(lc.SlotNum),
|
||||||
|
localcache.WithLocalSlotSize(lc.SlotSize),
|
||||||
|
localcache.WithLinkSlotNum(lc.SlotNum),
|
||||||
|
localcache.WithLocalSuccessTTL(lc.Success()),
|
||||||
|
localcache.WithLocalFailedTTL(lc.Failed()),
|
||||||
|
),
|
||||||
|
}
|
||||||
|
if lc.Enable() {
|
||||||
|
go subscriberRedisDeleteCache(context.Background(), cli, lc.Topic, x.local.DelLocal)
|
||||||
|
}
|
||||||
|
return x
|
||||||
|
}
|
||||||
|
|
||||||
|
type UserLocalCache struct {
|
||||||
|
client rpcclient.UserRpcClient
|
||||||
|
local localcache.Cache[any]
|
||||||
|
}
|
||||||
|
|
||||||
|
func (u *UserLocalCache) GetUserInfo(ctx context.Context, userID string) (val *sdkws.UserInfo, err error) {
|
||||||
|
log.ZDebug(ctx, "UserLocalCache GetUserInfo req", "userID", userID)
|
||||||
|
defer func() {
|
||||||
|
if err == nil {
|
||||||
|
log.ZDebug(ctx, "UserLocalCache GetUserInfo return", "value", val)
|
||||||
|
} else {
|
||||||
|
log.ZError(ctx, "UserLocalCache GetUserInfo return", err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
return localcache.AnyValue[*sdkws.UserInfo](u.local.Get(ctx, cachekey.GetUserInfoKey(userID), func(ctx context.Context) (any, error) {
|
||||||
|
log.ZDebug(ctx, "UserLocalCache GetUserInfo rpc", "userID", userID)
|
||||||
|
return u.client.GetUserInfo(ctx, userID)
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
|
func (u *UserLocalCache) GetUserGlobalMsgRecvOpt(ctx context.Context, userID string) (val int32, err error) {
|
||||||
|
log.ZDebug(ctx, "UserLocalCache GetUserGlobalMsgRecvOpt req", "userID", userID)
|
||||||
|
defer func() {
|
||||||
|
if err == nil {
|
||||||
|
log.ZDebug(ctx, "UserLocalCache GetUserGlobalMsgRecvOpt return", "value", val)
|
||||||
|
} else {
|
||||||
|
log.ZError(ctx, "UserLocalCache GetUserGlobalMsgRecvOpt return", err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
return localcache.AnyValue[int32](u.local.Get(ctx, cachekey.GetUserGlobalRecvMsgOptKey(userID), func(ctx context.Context) (any, error) {
|
||||||
|
log.ZDebug(ctx, "UserLocalCache GetUserGlobalMsgRecvOpt rpc", "userID", userID)
|
||||||
|
return u.client.GetUserGlobalMsgRecvOpt(ctx, userID)
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
|
func (u *UserLocalCache) GetUsersInfo(ctx context.Context, userIDs []string) ([]*sdkws.UserInfo, error) {
|
||||||
|
users := make([]*sdkws.UserInfo, 0, len(userIDs))
|
||||||
|
for _, userID := range userIDs {
|
||||||
|
user, err := u.GetUserInfo(ctx, userID)
|
||||||
|
if err != nil {
|
||||||
|
if errs.ErrRecordNotFound.Is(err) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
users = append(users, user)
|
||||||
|
}
|
||||||
|
return users, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (u *UserLocalCache) GetUsersInfoMap(ctx context.Context, userIDs []string) (map[string]*sdkws.UserInfo, error) {
|
||||||
|
users := make(map[string]*sdkws.UserInfo, len(userIDs))
|
||||||
|
for _, userID := range userIDs {
|
||||||
|
user, err := u.GetUserInfo(ctx, userID)
|
||||||
|
if err != nil {
|
||||||
|
if errs.ErrRecordNotFound.Is(err) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
users[userID] = user
|
||||||
|
}
|
||||||
|
return users, nil
|
||||||
|
}
|
Loading…
Reference in new issue