Compare commits

...

6 Commits

7 changed files with 87 additions and 31 deletions

View File

@ -143,14 +143,16 @@ func (x *JSSdk) getActiveConversations(ctx context.Context, req *jssdk.GetActive
if err != nil {
return nil, err
}
msgs, err := x.msgClient.GetSeqMessage(ctx, req.OwnerUserID, datautil.Slice(sortList, func(c *msg.ActiveConversation) *msg.ConversationSeqs {
return &msg.ConversationSeqs{
ConversationID: c.ConversationID,
Seqs: []int64{c.MaxSeq},
maxSeqs := datautil.SliceToMapAny(sortList, func(c *msg.ActiveConversation) (string, int64) {
return c.ConversationID, c.MaxSeq
})
conversationSeqs := x.filterConversationSeqs(conversations, maxSeqs)
var msgs map[string]*sdkws.PullMsgs
if len(conversationSeqs) > 0 {
msgs, err = x.msgClient.GetSeqMessage(ctx, req.OwnerUserID, conversationSeqs)
if err != nil {
return nil, err
}
}))
if err != nil {
return nil, err
}
x.checkMessagesAndGetLastMessage(ctx, req.OwnerUserID, msgs)
conversationMap := datautil.SliceToMap(conversations, func(c *conversation.Conversation) string {
@ -208,15 +210,7 @@ func (x *JSSdk) getConversations(ctx context.Context, req *jssdk.GetConversation
if err != nil {
return nil, err
}
conversationSeqs := make([]*msg.ConversationSeqs, 0, len(conversations))
for _, c := range conversations {
if seq := maxSeqs[c.ConversationID]; seq > 0 {
conversationSeqs = append(conversationSeqs, &msg.ConversationSeqs{
ConversationID: c.ConversationID,
Seqs: []int64{seq},
})
}
}
conversationSeqs := x.filterConversationSeqs(conversations, maxSeqs)
var msgs map[string]*sdkws.PullMsgs
if len(conversationSeqs) > 0 {
msgs, err = x.msgClient.GetSeqMessage(ctx, req.OwnerUserID, conversationSeqs)
@ -253,6 +247,21 @@ func (x *JSSdk) getConversations(ctx context.Context, req *jssdk.GetConversation
}, nil
}
func (x *JSSdk) filterConversationSeqs(conversations []*conversation.Conversation, maxSeqs map[string]int64) []*msg.ConversationSeqs {
conversationSeqs := make([]*msg.ConversationSeqs, 0, len(conversations))
for _, c := range conversations {
seq := maxSeqs[c.ConversationID]
if seq == 0 || (c.MinSeq > 0 && seq < c.MinSeq) {
continue
}
conversationSeqs = append(conversationSeqs, &msg.ConversationSeqs{
ConversationID: c.ConversationID,
Seqs: []int64{seq},
})
}
return conversationSeqs
}
// This function checks whether the latest MaxSeq message is valid.
// If not, it needs to fetch a valid message again.
func (x *JSSdk) checkMessagesAndGetLastMessage(ctx context.Context, userID string, messages map[string]*sdkws.PullMsgs) {

View File

@ -32,7 +32,7 @@ func TestCompressDecompress(t *testing.T) {
compressor := NewGzipCompressor()
for i := 0; i < 2000; i++ {
for range 2000 {
src := mockRandom()
// compress
@ -58,10 +58,8 @@ func TestCompressDecompressWithConcurrency(t *testing.T) {
wg := sync.WaitGroup{}
compressor := NewGzipCompressor()
for i := 0; i < 200; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for range 200 {
wg.Go(func() {
src := mockRandom()
// compress
@ -80,8 +78,7 @@ func TestCompressDecompressWithConcurrency(t *testing.T) {
// check
assert.EqualValues(t, src, res)
}()
})
}
wg.Wait()
}
@ -90,7 +87,7 @@ func BenchmarkCompress(b *testing.B) {
src := mockRandom()
compressor := NewGzipCompressor()
for i := 0; i < b.N; i++ {
for b.Loop() {
_, err := compressor.Compress(src)
assert.Equal(b, nil, err)
}
@ -100,7 +97,7 @@ func BenchmarkCompressWithSyncPool(b *testing.B) {
src := mockRandom()
compressor := NewGzipCompressor()
for i := 0; i < b.N; i++ {
for b.Loop() {
_, err := compressor.CompressWithPool(src)
assert.Equal(b, nil, err)
}
@ -114,7 +111,7 @@ func BenchmarkDecompress(b *testing.B) {
assert.Equal(b, nil, err)
for i := 0; i < b.N; i++ {
for b.Loop() {
_, err := compressor.DeCompress(comdata)
assert.Equal(b, nil, err)
}
@ -127,7 +124,7 @@ func BenchmarkDecompressWithSyncPool(b *testing.B) {
comdata, err := compressor.Compress(src)
assert.Equal(b, nil, err)
for i := 0; i < b.N; i++ {
for b.Loop() {
_, err := compressor.DecompressWithPool(comdata)
assert.Equal(b, nil, err)
}

View File

@ -1,9 +1,10 @@
package msggateway
import (
"github.com/openimsdk/tools/utils/datautil"
"sync"
"time"
"github.com/openimsdk/tools/utils/datautil"
)
type UserMap interface {
@ -146,7 +147,7 @@ func (u *userMap) DeleteClients(userID string, clients []*Client) (isDeleteUser
return client.ctx.GetRemoteAddr()
})
tmp := result.Clients
result.Clients = result.Clients[:0]
result.Clients = result.Clients[:0:0]
for _, client := range tmp {
if _, delCli := deleteAddr[client.ctx.GetRemoteAddr()]; delCli {
offline = append(offline, int32(client.PlatformID))

View File

@ -287,6 +287,19 @@ func (g *groupServer) webhookAfterJoinGroup(ctx context.Context, after *config.A
g.webhookClient.AsyncPost(ctx, cbReq.GetCallbackCommand(), cbReq, &callbackstruct.CallbackAfterJoinGroupResp{}, after)
}
func afterJoinGroupRequests(members []*model.GroupMember, reqMessage string) []*group.JoinGroupReq {
requests := make([]*group.JoinGroupReq, 0, len(members))
for _, member := range members {
requests = append(requests, &group.JoinGroupReq{
GroupID: member.GroupID,
ReqMessage: reqMessage,
JoinSource: member.JoinSource,
InviterUserID: member.UserID,
})
}
return requests
}
func (g *groupServer) webhookBeforeSetGroupInfo(ctx context.Context, before *config.BeforeConfig, req *group.SetGroupInfoReq) error {
return webhook.WithCondition(ctx, before, func(ctx context.Context) error {
cbReq := &callbackstruct.CallbackBeforeSetGroupInfoReq{

View File

@ -23,9 +23,10 @@ import (
"strings"
"time"
"github.com/openimsdk/tools/utils/stringutil"
"google.golang.org/grpc"
"github.com/openimsdk/tools/utils/stringutil"
"github.com/openimsdk/open-im-server/v3/pkg/dbbuild"
"github.com/openimsdk/open-im-server/v3/pkg/rpcli"
@ -477,6 +478,9 @@ func (g *groupServer) InviteUserToGroup(ctx context.Context, req *pbgroup.Invite
if err := g.setMemberJoinSeq(ctx, req.GroupID, req.InvitedUserIDs); err != nil {
return nil, err
}
for _, joinReq := range afterJoinGroupRequests(groupMembers, req.Reason) {
g.webhookAfterJoinGroup(ctx, &g.config.WebhooksConfig.AfterJoinGroup, joinReq)
}
return &pbgroup.InviteUserToGroupResp{}, nil
}
@ -913,6 +917,9 @@ func (g *groupServer) GroupApplicationResponse(ctx context.Context, req *pbgroup
if err := g.setMemberJoinSeq(ctx, req.GroupID, []string{req.FromUserID}); err != nil {
return nil, err
}
for _, joinReq := range afterJoinGroupRequests([]*model.GroupMember{member}, groupRequest.ReqMsg) {
g.webhookAfterJoinGroup(ctx, &g.config.WebhooksConfig.AfterJoinGroup, joinReq)
}
}
case constant.GroupResponseRefuse:
g.notification.GroupApplicationRejectedNotification(ctx, req)
@ -1315,6 +1322,9 @@ func (g *groupServer) TransferGroupOwner(ctx context.Context, req *pbgroup.Trans
}
func (g *groupServer) GetGroups(ctx context.Context, req *pbgroup.GetGroupsReq) (*pbgroup.GetGroupsResp, error) {
if err := authverify.CheckAdmin(ctx); err != nil {
return nil, err
}
var (
group []*model.Group
err error

View File

@ -22,6 +22,8 @@ import (
"sync"
"time"
"google.golang.org/grpc"
"github.com/openimsdk/open-im-server/v3/internal/rpc/relation"
"github.com/openimsdk/open-im-server/v3/pkg/authverify"
"github.com/openimsdk/open-im-server/v3/pkg/common/config"
@ -46,7 +48,6 @@ import (
"github.com/openimsdk/tools/discovery"
"github.com/openimsdk/tools/errs"
"github.com/openimsdk/tools/utils/datautil"
"google.golang.org/grpc"
)
const (
@ -268,6 +269,9 @@ func (s *userServer) AccountCheck(ctx context.Context, req *pbuser.AccountCheckR
}
func (s *userServer) GetPaginationUsers(ctx context.Context, req *pbuser.GetPaginationUsersReq) (resp *pbuser.GetPaginationUsersResp, err error) {
if err = authverify.CheckAdmin(ctx); err != nil {
return nil, err
}
if req.UserID == "" && req.NickName == "" {
total, users, err := s.db.PageFindUser(ctx, constant.IMOrdinaryUser, constant.AppOrdinaryUsers, req.Pagination)
if err != nil {
@ -353,6 +357,9 @@ func (s *userServer) GetGlobalRecvMessageOpt(ctx context.Context, req *pbuser.Ge
// GetAllUserID Get user account by page.
func (s *userServer) GetAllUserID(ctx context.Context, req *pbuser.GetAllUserIDReq) (resp *pbuser.GetAllUserIDResp, err error) {
if err = authverify.CheckAdmin(ctx); err != nil {
return nil, err
}
total, userIDs, err := s.db.GetAllUserID(ctx, req.Pagination)
if err != nil {
return nil, err

View File

@ -835,6 +835,25 @@ func (db *commonMsgDatabase) GetLastMessage(ctx context.Context, conversationIDs
}
return nil, err
}
if msg == nil || msg.Msg == nil {
continue
}
if userID != "" {
userMinSeq, err := db.seqUser.GetUserMinSeq(ctx, conversationID, userID)
if err != nil {
return nil, err
}
minSeq, err := db.seqConversation.GetMinSeq(ctx, conversationID)
if err != nil {
return nil, err
}
if userMinSeq > minSeq {
minSeq = userMinSeq
}
if msg.Msg.Seq < minSeq {
continue
}
}
tmp := []*model.MsgInfoModel{msg}
db.handlerDeleteAndRevoked(ctx, userID, tmp)
db.handlerQuote(ctx, userID, conversationID, tmp)